拓冰建站拓冰建站
首页 / 资讯中心 / 正文

SeaTunnel 作业环境(env)配置完整指南:通用参数、Zeta 专属参数与引擎前缀规则

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 支持在 ZetaSeaTunnel 自研引擎、Flink、Spark 上运行同一份任务配置。为了区分参数归属官方制定了明确的前缀规则通用参数不携带任何前缀可在所有引擎中使用Flink 引擎参数必须携带flink.前缀例如flink.pipeline.max-parallelismSpark 引擎参数不添加任何额外前缀。原因是 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还是流处理STREAMINGenv { 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.yamlseatunnel: 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作业失败重试次数控制作业失败时的默认重试次数默认值为 3env { 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.locationSaveMode 执行位置指定作业在 Zeta 引擎下执行 SaveMode保存模式即写入目标表前的建表/清表等预处理动作的位置env { savemode.execute.location CLUSTER # 或 savemode.execute.location CLIENT }默认值为CLUSTERSaveMode 在集群端执行若需在客户端执行可设置为CLIENT官方强烈建议使用CLUSTER模式文档明确说明当CLUSTER模式不再存在问题时CLIENT模式将被移除。源码中该参数是枚举类型默认值为SaveModeExecuteLocation.CLUSTER见 EnvCommonOptions.java在 Zeta 引擎的 MultipleTableJobConfigParser.java 等配置解析链路中被读取并分发到对应执行位置。4.4 sink.flush.intervalSink 主动冲刷间隔引擎向流水线注入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-parallelismflink.pipeline.max-parallelismexecution.checkpointing.modeflink.execution.checkpointing.modeexecution.checkpointing.timeoutflink.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 批处理任务BATCHenv { job.name daily_batch_sync job.mode BATCH parallelism 4 # BATCH 模式不设置 checkpoint.interval 即禁用检查点 }7.2 流处理任务STREAMINGZeta 引擎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.nameStringSeaTunnel_Job全部任务名称jarsString无全部第三方依赖包分号分隔job.modeEnumBATCH全部批/流模式必填checkpoint.intervalLong无STREAMING 下 Zeta 兜底 30000ms全部检查点间隔毫秒checkpoint.timeoutLong无Zeta 默认 30000ms全部检查点超时毫秒parallelismInteger1全部源/汇默认并行度shade.identifierString无全部配置文件加解密方式job.retry.timesInteger3仅 Zeta失败重试次数生命周期内累计job.retry.interval.secondsInteger3仅 Zeta失败重试间隔秒savemode.execute.locationEnumCLUSTER仅 ZetaSaveMode 执行位置CLUSTER/CLIENTsink.flush.intervalLong0禁用仅 ZetaSink 冲刷信号注入间隔毫秒九、源码定位与进一步阅读所有 env 参数的声明、类型与默认值EnvCommonOptions.javaenv 参数的必填/可选规则EnvOptionRuleFactoryEnvOptionRule.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),仅供参考
分享:

看完干货,该让你的企业上线了

免费需求沟通 · 48 小时内出具建站方案 · 河南本地可上门