SeaTunnel Python 源连接器(Python Source Connector)实战指南:脚本拉起、stdin/stdout 协议与安全白名单
【免费下载链接】seatunnelSeaTunnel is a multimodal, high-performance, distributed, massive data integration tool.项目地址: https://gitcode.com/GitHub_Trending/se/seatunnel
本文围绕 Apache SeaTunnel 的 Python 源连接器(插件名
Python)展开,讲解其「拉起 Python 脚本、以 stdin 第一行传入 JSON 配置、逐行解析 stdout 为 SeaTunnel Row」的核心设计、全部配置选项、集群侧安全白名单机制与完整可运行的实战示例。读完本文,你将掌握如何在 SeaTunnel 作业中接入任意 Python 脚本作为数据源,并理解其 Phase 1 MVP 的实现边界与源码级行为。
概述:为什么需要 Python 源连接器
在真实的数据集成场景中,很多数据源没有现成的 JDBC、HTTP 或消息队列驱动,但业务方往往已经沉淀了成熟的 Python 取数逻辑(调用内部 API、解析私有协议、访问加密存储等)。SeaTunnel 的 Python 源连接器正是为了复用这类存量 Python 能力而设计:它只要求脚本「从 stdin 读一行 JSON、向 stdout 打印文本行」,SeaTunnel 侧负责进程生命周期管理与行解析,Python 侧只需要专注取数与打印,实现门槛很低。
从源码看,该连接器位于 seatunnel-connectors-v2/connector-python 模块,并在 config/plugin_config 中登记为connector-python插件,可通过标准的连接器发现机制加载(工厂类使用@AutoService注册,见 PythonSourceFactory.java)。
工作原理:一次进程启动,两条 I/O 约定
Phase 1 MVP 的实现模型非常精简,可用一句话概括:SeaTunnel 通过ProcessBuilder启动一个 Python 进程,把python.script.config序列化成 JSON 写到该进程 stdin 的第一行,再把 stdout 的每一行按 text 格式解析成 SeaTunnel Row。
对应到源码,PythonSourceReader.java 中进程启动的核心逻辑是:
ProcessBuilder processBuilder = new ProcessBuilder(resolvedExecutable.toString(), scriptPath.toString()); configureWorkingDirectory(processBuilder, scriptPath);即启动命令等价于python3 <script_path>,工作目录默认取脚本所在目录。
stdin/stdout 交互约定(重要)
:::tip Python 交互约定
当前 MVP 只支持file_format_type = text。
- SeaTunnel 会把
python.script.config作为 JSON 写到 stdin 第一行。 - Python 脚本每打印一行 stdout,就表示一条记录。
- SeaTunnel 会按
field_delimiter切分该行,并依据schema做类型转换。
:::
这一约定的实现证据在 PythonSourceReader.java:配置以 UTF-8 编码写入一行 JSON 并换行 flush;stdout 侧则由TextDeserializationSchema(见 createDeserializationSchema)按field_delimiter分隔并按schema声明的类型逐字段转换。
关键特性一览
| 特性 | 支持情况 |
|---|---|
| 批(Batch) | ✅ 支持 |
| 流(Stream) | ❌ 不支持 |
| 精确一次(Exactly-once) | ❌ 不支持 |
| 列投影(Column Projection) | ✅ 支持 |
| 并行性(Parallelism) | ❌ 不支持 |
| 支持用户自定义 split | ❌ 不支持 |
上述特性语义可参考 Connector V2 功能简介。源码层面,PythonSource.java 将数据源声明为Boundedness.BOUNDED(有界),并实现了SupportColumnProjection接口以支持列投影。
选项详解
| 参数名 | 类型 | 必须 | 默认值 |
|---|---|---|---|
| python.executable | string | 否 | python3 |
| python.script.path | string | 是 | - |
| schema | config | 是 | - |
| python.script.config | map | 否 | {} |
| python.working.directory | string | 否 | 脚本所在目录 |
| file_format_type | string | 否 | text |
| field_delimiter | string | 否 | , |
| common-options | 否 | - |
必填项与可选项的定义在 PythonSourceFactory.java 的OptionRule中声明:python.script.path与schema必填,其余为可选。各选项的默认值、类型与描述均可从 PythonSourceOptions.java 中逐一定义确认。
python.executable [string]
用于启动脚本的 Python 解释器或可执行文件。默认值为python3。
- 可以填写绝对路径(如
/usr/bin/python3、/opt/venv/bin/python)或裸命令名(如python3),不允许填写带路径前缀的相对形式; - 若填写裸命令名,SeaTunnel 会按
PATH环境变量解析出绝对路径; - 最终解析出的绝对路径必须包含在集群管理员控制的系统属性
seatunnel.source.python.allowed-executables白名单中,否则进程启动会被拒绝。
解析与白名单校验逻辑位于 PythonSourceExecutionPolicy.java,详见下文「安全模型」。
python.script.path [string]
要执行的 Python 脚本路径,必填。SeaTunnel 在启动进程前会校验该路径必须指向一个真实存在的常规文件(Files.isRegularFile),否则直接抛出IllegalArgumentException快速失败,见 validateScriptPath。
schema [config]
stdout 记录的 schema,必填。SeaTunnel 会按这个 schema 把每一行 text 输出转换成SeaTunnelRow。声明方式支持fields简写形式:
schema = { fields { id = int name = string } }也支持完整的columns、primaryKey、constraintKeys等高级定义。支持的字段类型包括string、boolean、tinyint、smallint、int、bigint、float、double、decimal、date、time、timestamp以及row/map/array等复合类型,详细语法请参考 Schema 特性简介。
python.script.config [map]
可选配置对象,默认{}。SeaTunnel 会把它序列化成 JSON,写入 Python 进程 stdin 的第一行。
它适合传递 API 地址、鉴权信息、过滤条件或其他运行时参数,避免把这些配置硬编码进脚本。例如:
python.script.config = { prefix = "seatunnel" count = 3 }脚本侧通过json.loads(sys.stdin.readline())读取即可。源码中该 JSON 的序列化使用JsonUtils.toJsonString(见 writeInitialScriptConfig)。
python.working.directory [string]
Python 进程的工作目录。未配置时,默认使用python.script.path的父目录。配置时该目录必须真实存在(Files.isDirectory校验),实现见 configureWorkingDirectory。
file_format_type [string]
stdout 解析格式,默认text。当前 Phase 1 只支持:
text
若配置其他值,PythonSourceConfig的校验逻辑会直接抛出Unsupported file_format_type ... Phase 1 supports only text,实现见 PythonSourceConfig.java。
field_delimiter [string]
当file_format_type = text时使用的字段分隔符,默认,。示例:,、|、\t。该分隔符会传给TextDeserializationSchema用于切分每一行 stdout。
common options
源插件通用参数,包括plugin_output、parallelism、metadata_datasource_id等,请参考 Source 常用选项。其中plugin_output用于把本插件产出的数据注册为可被下游通过plugin_input直接访问的临时表/数据集。
安全模型:默认禁用与双重白名单
Python 源连接器会以 SeaTunnel worker 进程的权限在 Worker 节点上直接执行外部代码,因此安全策略是使用前必须理解的重中之重。
- 该连接器默认禁用。集群管理员必须在每个 Worker 节点设置 JVM 系统属性:
-Dseatunnel.source.python.enabled=true(总开关)-Dseatunnel.source.python.allowed-executables=/absolute/path/to/python3(可执行文件白名单,多个路径用逗号分隔,且必须是绝对路径)
- 任务配置不能启用该连接器,也不能扩大该白名单。也就是说,开关与白名单只能由集群侧系统属性控制,作业配置无法越权。
python.executable和python.script.path会在 worker 节点上以 SeaTunnel worker 进程的权限直接执行。- 每次启动进程时,SeaTunnel 都会把最终解析出的可执行文件和规范化后的
python.script.path作为审计告警写入日志(以LOG.warn级别输出,见 PythonSourceReader.open)。 python.script.config会被序列化为 JSON 并写入子进程 stdin,因此其中的密钥或令牌会暴露给该子进程及其运行期日志或诊断信息,请勿放置明文敏感凭据。- 在共享集群中,建议限制谁可以提交使用该连接器的任务,并尽量让 worker 运行在受控或隔离的环境里。
白名单校验的源码逻辑
PythonSourceExecutionPolicy.java 完整实现了三层校验:
ensureEnabled():读取seatunnel.source.python.enabled,未显式置true即抛异常;parseAllowedExecutables():解析seatunnel.source.python.allowed-executables,要求至少一个绝对路径条目,且路径规范化后去重;resolveConfiguredExecutable()+ 白名单比对:把作业配置的python.executable解析为绝对路径(裸命令名走PATH解析),再与白名单逐项比对(先比较规范化路径,再用Files.isSameFile兜底处理软链接等场景),不在白名单内即拒绝启动。
这套校验同时被单测覆盖:PythonSourceTest中testReaderRejectsPythonExecutionWhenServerPolicyIsDisabled与testReaderRejectsExecutableOutsideServerAllowlist分别验证了「未开启开关」和「可执行文件不在白名单」两种拒绝场景(见 PythonSourceTest.java)。
完整实战示例
下面给出一个端到端可运行的示例:Python 脚本从 stdin 读取配置,打印count行由prefix拼接的记录,SeaTunnel 将其解析为(id int, name string)两列并输出到 Console。
SeaTunnel 配置
env { parallelism = 1 job.mode = "BATCH" } source { Python { plugin_output = "python_source" python.executable = "/usr/bin/python3" python.script.path = "/tmp/python_source.py" python.script.config = { prefix = "seatunnel" count = 3 } file_format_type = "text" field_delimiter = "," schema = { fields { id = int name = string } } } } sink { Console { plugin_input = "python_source" } }Python 脚本
#!/usr/bin/env python3 import json import sys def main(): config_line = sys.stdin.readline().strip() config = json.loads(config_line) if config_line else {} prefix = config.get("prefix", "python") count = int(config.get("count", 2)) for index in range(1, count + 1): print(f"{index},{prefix}_{index}", flush=True) if __name__ == "__main__": main()运行与验证要点
- 上述脚本与仓库测试资源 emit_rows.py 的结构一致(该脚本在单测中用于验证整条数据链路)。
- 作业输出预期为 3 行数据:
(1, seatunnel_1)、(2, seatunnel_2)、(3, seatunnel_3)。 - 注意脚本中
print(..., flush=True)是必须的:只有主动 flush,stdout 才能及时被 SeaTunnel 的 stdout 管道读取。 - 提交作业前请确保 Worker 节点已配置
-Dseatunnel.source.python.enabled=true且白名单包含/usr/bin/python3,否则连接器会直接拒绝启动并抛出异常。
源码级实现剖析
连接器类结构
connector-python模块的源码结构非常清晰,共 6 个类,全部位于 source 包:
| 类 | 职责 |
|---|---|
| PythonSource | Source 主类,声明有界性(BOUNDED)、插件名与列投影支持 |
| PythonSourceFactory | 工厂类,负责插件注册与 OptionRule(必填/可选参数)声明 |
| PythonSourceOptions | 全部配置项的 Option 定义(key、类型、默认值、描述) |
| PythonSourceConfig | 运行时配置快照,构造时完成校验(快速失败) |
| PythonSourceReader | Reader 实现,管理进程生命周期、三条 pump 线程与行解析 |
| PythonSourceExecutionPolicy | 安全策略:开关检查、可执行文件解析与白名单比对 |
工厂类的参数契约同样有单测保护:PythonSourceFactoryTest.java 验证了工厂标识符、必填/可选参数集合以及参数命名。
Reader 生命周期与并发控制
PythonSourceReader.java 是连接器最核心的实现,其关键设计包括:
- 三条后台线程:
stdin-writer(写 JSON 配置并关闭 stdin)、stdout-pump(逐行读取 stdout 放入有界队列)、stderr-pump(把 stderr 转发到 worker 日志并保留最近 50 行用于错误上下文),全部为 daemon 线程; - 有界队列与背压:stdout 行放入容量为 256 的
ArrayBlockingQueue,每次pollNext最多发射 128 行,从源码看这是通过STDOUT_QUEUE_CAPACITY = 256与MAX_ROWS_PER_POLL = 128两个常量控制; - 进程完成检测:进程退出后先等待 stdout EOF,把缓冲行全部排空后才调用
signalNoMoreElement()通知引擎有界数据源结束;若脚本派生的子进程继承了 stdout 管道导致迟迟不关闭,则超过 5 秒宽限期后会显式抛出 "ensure child processes do not inherit stdout" 的协议错误,而不是无限挂起; - 非零退出处理:若 Python 进程以非零码退出,Source task 会失败,异常信息中携带最近 50 行 stderr 输出(如
exited with code 1. Recent stderr: ...); - 关闭语义:
close()先destroy(),等待 5 秒未退出则destroyForcibly(),且通过lifecycleLock协调并发关闭与正在进行的 poll,避免竞态。
测试用例如何验证这些行为
PythonSourceTest.java 提供了非常完整的单元测试矩阵,可以直接作为理解连接器行为的「可执行文档」:
testReaderCollectsRowsFromPythonScript:验证整条链路(启动进程 → 传配置 → 解析 stdout → 收集 2 行并通知完成);testReaderDrainsBufferedRowsAfterProcessExit:验证输出超过有界队列容量(300 行)时,进程退出后缓冲行仍能完整排空;testReaderCollectsTrailingRowWithoutNewline:验证最后一行无换行符时也能在 EOF 处被捕获;testReaderFailsWhenPythonProcessExitsNonZero:验证非零退出时异常携带 stderr 内容;testReaderCloseStopsLongRunningPythonProcess:验证 close 能在 5 秒内终止无限运行的脚本;testReaderOpenTimesOutWhenPythonDoesNotReadLargeConfig:验证脚本不读 stdin 时,写配置在 5 秒超时后失败;testReaderFailsWhenChildKeepsStdoutOpen:验证子进程继承 stdout 时会显式失败而非挂起。
限制与注意事项
Phase 1 MVP 的使用边界(务必在选型与架构设计时考虑):
- 只支持 source,尚不支持作为 sink 使用;
- 只支持
text输出格式,且每次脚本执行产生一条有限 stdout 流,因此被建模为有界(BATCH)单 split 数据源; - 当前 Source 只有单 reader,source parallelism 必须保持为
1; - 不保存任何可恢复的位点或 checkpoint 状态。任务失败恢复或重启后,Python 脚本会从头重新执行,之前已经下发的行会再次发出;请使用幂等的 sink,或确保任务可以容忍重复数据;
- 连接器只管理直接启动的进程。脚本不能派生继承 stdout 或 stderr 的长期后台子进程;如需管理子进程树,应由 worker 侧的隔离与进程监管机制负责;
- 如果 Python 进程非零退出,Source task 会失败,并在异常里带上最近的 stderr 输出(便于快速定位脚本问题)。
变更日志
[Feature][Connector-V2] Add Python source connector (#11337)— 新增 Python 源连接器,版本标记为 Next,详见 connector-python 变更日志。
【免费下载链接】seatunnelSeaTunnel is a multimodal, high-performance, distributed, massive data integration tool.项目地址: https://gitcode.com/GitHub_Trending/se/seatunnel
创作声明:本文部分内容由AI辅助生成(AIGC),仅供参考