SeaTunnel 作业环境(env)配置完整指南:通用参数、Zeta 专属参数与引擎前缀规则
【免费下载链接】seatunnelSeaTunnel is a multimodal, high-performance, distributed, massive data integration tool.项目地址: https://gitcode.com/GitHub_Trending/se/seatunnel
导读:
env配置块是 SeaTunnel 任务配置文件(HOCON 格式)中的首个核心块,用于声明任务名称、运行模式(批/流)、并行度、检查点行为、重试策略等作业级参数。本文以官方文档 JobEnvConfig.md 为主体,结合 EnvCommonOptions.java 等源码,系统讲解每个参数的类型、默认值、生效引擎及底层实现,并给出可直接复制运行的完整配置示例。读完本文,你将掌握如何为批处理、流处理、多引擎迁移场景编写正确且可维护的env配置。
一、env 配置块:任务配置文件的起点
在 SeaTunnel 的 V2 配置体系中,一个完整的任务配置由env、source、transform(可选)、sink四个块组成。其中env块负责作业级(Job-level)环境参数,例如官方模板 config/v2.batch.config.template 中的写法:
env { # You can set SeaTunnel environment configuration here parallelism = 2 job.mode = "BATCH" checkpoint.interval = 10000 } source { FakeSource { parallelism = 2 plugin_output = "fake" row.num = 16 schema = { fields { name = "string" age = "int" } } } } sink { Console { } }env块中的每个键都是 SeaTunnel 定义的标准参数,它们在 EnvCommonOptions.java 中统一声明,并由 EnvOptionRule.java(一个基于OptionRule的Factory实现)约束哪些参数必填、哪些可选。从该源码可以看到:job.mode是唯一必填的 env 参数,其余(job.name、parallelism、job.retry.times、jars、checkpoint.interval、checkpoint.timeout、savemode.execute.location、sink.flush.interval等)均为可选。
二、引擎参数前缀规则:通用参数与引擎专属参数
SeaTunnel 支持在 Zeta(SeaTunnel 自研引擎)、Flink、Spark 上运行同一份任务配置。为了区分参数归属,官方制定了明确的前缀规则:
- 通用参数:不携带任何前缀,可在所有引擎中使用;
- Flink 引擎参数:必须携带
flink.前缀,例如flink.pipeline.max-parallelism; - Spark 引擎参数:不添加任何额外前缀。原因是 Spark 官方参数本身就以
spark.开头(如spark.executor.memory),如果 SeaTunnel 再加一层前缀会造成冲突,因此直接透传原生参数。
Zeta 引擎专属参数(如job.retry.times、job.retry.interval.seconds、savemode.execute.location、sink.flush.interval)则不带前缀,仅在 Zeta 引擎下生效。
三、通用参数详解(所有引擎生效)
以下参数在 EnvCommonOptions.java 中被定义为跨引擎通用选项,任何引擎下均可使用。
3.1 job.name:任务名称
配置任务的显示名称,用于在引擎的 Job 列表、Web UI、日志中标识该任务。
env { job.name = "SeaTunnel_Job" }从源码实现看,其键为job.name,类型为字符串,默认值为SeaTunnel_Job(见 EnvCommonOptions.java)。在 Zeta 引擎下,该名称会直接用于集群中的作业注册与 REST API 查询。
3.2 jars:加载第三方依赖包
通过jars可以加载第三方依赖包,多个 jar 使用分号;分隔:
env { jars = "file://local/jar1.jar;file://local/jar2.jar" }源码中该参数为字符串类型且无默认值(见 EnvCommonOptions.java)。在 Zeta 引擎的测试资源中可以看到实际用法,例如client_test_with_jars.conf(位于 seatunnel-engine/seatunnel-engine-client/src/test/resources),用于验证作业提交时附带第三方 jar 的加载链路。
注意:
jars通常用于连接器未内置、需要在任务级临时补充的依赖。若依赖属于某个连接器的固定依赖,更推荐将其加入该连接器的 lib 目录或通过plugin_config管理,避免每个任务重复声明。
3.3 job.mode:批处理 / 流处理模式
通过job.mode声明任务是批处理(BATCH)还是流处理(STREAMING):
env { job.mode = "BATCH" # 或 job.mode = "STREAMING" }源码层面,该参数是枚举类型,取值对应 JobMode.java 中的BATCH和STREAMING两个枚举常量,默认值为BATCH(见 EnvCommonOptions.java)。同时它是 EnvOptionRule.java 中唯一被标记为required的 env 参数,也就是说每个任务都必须显式声明job.mode。
job.mode直接决定检查点(Checkpoint)的行为:STREAMING模式下检查点是必须启用的;BATCH模式下可以不配置检查点参数来禁用检查点。
3.4 checkpoint.interval:检查点触发间隔
配置周期性调度检查点的时间间隔,单位是毫秒:
env { job.mode = "STREAMING" checkpoint.interval = 60000 # 每 60 秒触发一次检查点 }参数行为因模式而异:
- STREAMING 模式:检查点是必需的。如果未在任务配置中设置
checkpoint.interval,会从引擎的应用配置文件seatunnel.yaml中获取;在 Zeta 引擎的 STREAMING 模式下,默认值为30000 毫秒(30 秒)。 - BATCH 模式:可以不设置该参数,此时检查点被禁用。
引擎级兜底值来自seatunnel.yaml中的seatunnel.engine.checkpoint.interval,见官方默认配置 config/seatunnel.yaml:
seatunnel: engine: checkpoint: interval: 10000 timeout: 60000源码中该参数键为checkpoint.interval,类型为Long,无默认值(noDefaultValue(),见 EnvCommonOptions.java),这一设计与"STREAMING 下从 seatunnel.yaml 兜底、BATCH 下可禁用"的文档描述一致。
3.5 checkpoint.timeout:检查点超时时间
检查点超时时间(毫秒)。如果检查点在超时前未完成,任务将失败:
env { checkpoint.timeout = 60000 # 检查点 60 秒内必须完成 }在 Zeta 引擎中,默认值为30000 毫秒(30 秒)。与checkpoint.interval相同,源码中该参数类型为Long、无默认值(见 EnvCommonOptions.java),并可通过seatunnel.yaml的seatunnel.engine.checkpoint.timeout设置集群级默认值。
3.6 parallelism:源与汇的并行度
配置 source 与 sink 的并行度:
env { parallelism = 4 }源码中的语义值得注意(见 EnvCommonOptions.java):当连接器未显式指定parallelism时,才使用 env 中的并行度作为默认值;如果连接器自身指定了并行度,则以连接器的配置为准。默认值为 1。例如官方模板中 env 与 FakeSource 都写了parallelism,实际生效的是连接器级(FakeSource 的parallelism = 2)配置。
3.7 shade.identifier:配置文件加解密方式
指定配置文件加密/解密的方式。如果没有配置文件加密需求,可以忽略该参数:
env { # 例如使用 Base64 编码方式(具体取值见加密解密文档) # shade.identifier = "base64" }该参数与 SeaTunnel 的 Config Encryption/Decryption 功能配套使用,完整说明请参见 Config Encryption Decryption。仓库中对应的测试用例 RestApiSubmitJobConfigShadeDecryptTest.java 验证了 REST API 提交加密配置时解密链路的正确性。
四、Zeta 引擎专属参数
以下参数仅在 SeaTunnel Zeta 引擎下生效。
4.1 job.retry.times:作业失败重试次数
控制作业失败时的默认重试次数,默认值为 3:
env { job.retry.times = 5 }关键语义(官方文档明确说明):
- 重试计数器在流水线(pipeline)整个生命周期内累积,中途成功恢复不会重置计数器;
- 例如设置
job.retry.times = 5:流水线失败后重试,并在第 3 次尝试时恢复;此后再次失败,则只剩 2 次重试机会(第 4、5 次尝试),预算不会刷新回 5; - 唯一例外:Zeta 集群中发生 active-master 故障转移(failover)时,会从头重建流水线执行计划(及其重试计数器)。
源码中该参数键为job.retry.times、类型Integer、默认值 3(见 EnvCommonOptions.java),与文档描述完全一致。
4.2 job.retry.interval.seconds:失败重试间隔
控制作业失败后的重试间隔,默认值为 3 秒:
env { job.retry.interval.seconds = 10 }源码中该参数键为job.retry.interval.seconds、类型Integer、默认值 3(见 EnvCommonOptions.java)。
4.3 savemode.execute.location:SaveMode 执行位置
指定作业在 Zeta 引擎下执行 SaveMode(保存模式,即写入目标表前的建表/清表等预处理动作)的位置:
env { savemode.execute.location = "CLUSTER" # 或 savemode.execute.location = "CLIENT" }- 默认值为
CLUSTER:SaveMode 在集群端执行; - 若需在客户端执行,可设置为
CLIENT; - 官方强烈建议使用
CLUSTER模式:文档明确说明,当CLUSTER模式不再存在问题时,CLIENT模式将被移除。
源码中该参数是枚举类型,默认值为SaveModeExecuteLocation.CLUSTER(见 EnvCommonOptions.java),在 Zeta 引擎的 MultipleTableJobConfigParser.java 等配置解析链路中被读取并分发到对应执行位置。
4.4 sink.flush.interval:Sink 主动冲刷间隔
引擎向流水线注入FlushSignal(冲刷信号)以驱动 Sink 执行 flush 的间隔,单位毫秒。0或不设置(默认)表示禁用该机制,仅 Zeta 引擎生效:
env { sink.flush.interval = 5000 # 每 5 秒向 Sink 注入一次冲刷信号 }官方对取值的建议:
- 不建议设置低于 100ms 的值:过密的冲刷信号会占用流水线队列容量,挤占正常数据记录,且在尚无数据缓冲时触发空冲刷(empty flush),徒增 Sink 的 I/O 开销。
源码中该参数键为sink.flush.interval、类型Long、默认值 0L,且源码注释同样注明"低于 100ms 的值会记录 WARN 日志"(见 EnvCommonOptions.java),与文档建议相互印证。
五、Flink 引擎参数映射
在 Flink 引擎下,SeaTunnel 参数与 Flink 原生配置通过flink.前缀进行映射。以下是官方文档给出的部分映射关系(并非全部,完整列表请以 Flink 官方文档为准):
| Flink 配置名称 | SeaTunnel 配置名称 |
|---|---|
pipeline.max-parallelism | flink.pipeline.max-parallelism |
execution.checkpointing.mode | flink.execution.checkpointing.mode |
execution.checkpointing.timeout | flink.execution.checkpointing.timeout |
... | ... |
配置示例:
env { job.mode = "STREAMING" flink.execution.checkpointing.mode = "EXACTLY_ONCE" flink.pipeline.max-parallelism = 16 }六、Spark 引擎参数
由于 Spark 的配置项本身没有被 SeaTunnel 改写(透传原生配置),官方文档未在 SeaTunnel 文档中逐一罗列 Spark 参数清单,需要查阅 Spark 官方文档。Spark 参数无需额外前缀,直接书写即可:
env { spark.executor.memory = "2g" spark.executor.cores = 2 }七、典型完整配置示例
7.1 批处理任务(BATCH)
env { job.name = "daily_batch_sync" job.mode = "BATCH" parallelism = 4 # BATCH 模式不设置 checkpoint.interval 即禁用检查点 }7.2 流处理任务(STREAMING,Zeta 引擎)
env { job.name = "realtime_sync" job.mode = "STREAMING" checkpoint.interval = 60000 # 每 60 秒触发检查点 checkpoint.timeout = 30000 # 检查点 30 秒超时 job.retry.times = 5 # 失败重试 5 次(生命周期内累计) job.retry.interval.seconds = 10 # 重试间隔 10 秒 sink.flush.interval = 5000 # 每 5 秒驱动一次 Sink flush }7.3 加载第三方 jar 的任务
env { job.mode = "BATCH" jars = "file:///opt/lib/mysql-connector.jar;file:///opt/lib/udf.jar" }八、参数速查表
| 参数名 | 类型 | 默认值 | 生效引擎 | 说明 |
|---|---|---|---|---|
job.name | String | SeaTunnel_Job | 全部 | 任务名称 |
jars | String | 无 | 全部 | 第三方依赖包,分号分隔 |
job.mode | Enum | BATCH | 全部 | 批/流模式,必填 |
checkpoint.interval | Long | 无(STREAMING 下 Zeta 兜底 30000ms) | 全部 | 检查点间隔(毫秒) |
checkpoint.timeout | Long | 无(Zeta 默认 30000ms) | 全部 | 检查点超时(毫秒) |
parallelism | Integer | 1 | 全部 | 源/汇默认并行度 |
shade.identifier | String | 无 | 全部 | 配置文件加解密方式 |
job.retry.times | Integer | 3 | 仅 Zeta | 失败重试次数(生命周期内累计) |
job.retry.interval.seconds | Integer | 3 | 仅 Zeta | 失败重试间隔(秒) |
savemode.execute.location | Enum | CLUSTER | 仅 Zeta | SaveMode 执行位置(CLUSTER/CLIENT) |
sink.flush.interval | Long | 0(禁用) | 仅 Zeta | Sink 冲刷信号注入间隔(毫秒) |
九、源码定位与进一步阅读
- 所有 env 参数的声明、类型与默认值:EnvCommonOptions.java
- env 参数的必填/可选规则(
EnvOptionRuleFactory):EnvOptionRule.java - 批/流模式枚举定义:JobMode.java
- 引擎级检查点兜底配置(
seatunnel.engine.checkpoint.*):config/seatunnel.yaml - 完整 env 配置示例模板:config/v2.batch.config.template、config/v2.streaming.conf.template
- 配置加密解密(
shade.identifier)专项文档:Config Encryption Decryption - 带
jars的作业提交测试用例:client_test_with_jars.conf - 加密配置解密链路测试:RestApiSubmitJobConfigShadeDecryptTest.java
最后提醒:本文中标注"Zeta 默认值"的参数(如
checkpoint.interval、checkpoint.timeout的 30000ms)均指 Zeta 引擎行为;迁移到 Flink/Spark 引擎时,请以对应引擎官方文档的参数语义为准,并注意flink.前缀规则与 Spark 无前缀规则。
【免费下载链接】seatunnelSeaTunnel is a multimodal, high-performance, distributed, massive data integration tool.项目地址: https://gitcode.com/GitHub_Trending/se/seatunnel
创作声明:本文部分内容由AI辅助生成(AIGC),仅供参考