SeaTunnel 引擎(Zeta)本地快速开始:单机跑通第一个批处理作业
2026/9/20 7:24:15 网站建设 项目流程

SeaTunnel 引擎(Zeta)本地快速开始:单机跑通第一个批处理作业

【免费下载链接】seatunnelSeaTunnel is a multimodal, high-performance, distributed, massive data integration tool.项目地址: https://gitcode.com/GitHub_Trending/se/seatunnel

SeaTunnel 引擎(Zeta)是 SeaTunnel 的默认内置引擎,既可以在一台机器上以 Local 模式快速验证配置、连接器与处理链路,也可以部署为多节点集群承接测试、预发和生产任务。本文以 SeaTunnel 引擎快速开始 为主线,完整演示从下载部署、安装插件、编写 HOCON 作业配置到运行与结果验证的全过程,并额外给出「MySQL 到 Doris」的真实批处理示例与 Local 模式的底层行为说明,帮助你一次性跑通首个 SeaTunnel 引擎作业。

两种使用方式与适用场景

SeaTunnel Engine 支持两种组织方式,本文对应表格中的两种路径:

使用方式适用场景下一步
单机快速开始在一台机器上验证配置、连接器或处理链路继续阅读本文的单机快速开始部分
集群部署在测试、预发或生产环境中运行多节点任务跳转到 SeaTunnel Engine(Zeta) 安装部署

从引擎部署模式看(参见 deployment.md),Zeta 支持本地模式、混合集群模式和分离集群模式三种形态。Local 模式只用于测试,每个任务都会启动一个独立进程,任务运行完成后进程退出;混合集群模式中 Master 与 Worker 同进程且所有节点可参与选举;分离集群模式则将 Master 服务与 Worker 服务拆分为独立进程,是官方建议的生产部署形态。

开始前建议先看

如果你是第一次接触 SeaTunnel 文档,建议按下列顺序建立整体路径感:

  • 快速入门总览
  • 安装部署
  • 作业配置指南

本文示例链路使用FakeSourceFieldMapperConsole三个插件,不依赖任何外部中间件,适合在单机上完成端到端验证。

第一部分:单机快速开始(Local 模式)

Local 模式适合在单台机器上快速验证安装、连接器和作业配置,下面的命令统一使用-m local启动 SeaTunnel Engine。启动时系统会在提交作业的进程中直接拉起引擎服务来运行作业,作业完成后进程随之退出,无需预部署任何集群组件。

步骤 1:部署 SeaTunnel 及连接器

在开始前,请确保已按照 部署 中的描述下载并部署 SeaTunnel:

  1. 安装 Java 8 或 11(其他高于 Java 8 的版本理论上也可以工作)并设置JAVA_HOME
  2. 下载二进制安装包seatunnel-<version>-bin.tar.gz并解压;
  3. 从 2.2.0-beta 版本开始,二进制包不再默认提供连接器依赖,首次使用需要手动安装插件。

如果你已经安装了完整插件,可以直接复用;若只是为了以最短路径跑通本文示例,只需connector-fakeconnector-console两个插件即可。

步骤 2:安装示例所需插件

编辑${SEATUNNEL_HOME}/config/plugin_config,只保留示例需要的两个插件:

--seatunnel-connectors-- connector-fake connector-console --end--

然后执行安装脚本(脚本会读取plugin_config并从仓库拉取对应 JAR 到${SEATUNNEL_HOME}/connectors/目录):

sh bin/install-plugin.sh

所有支持的连接器及其在plugin_config中对应的配置名称,可以在${SEATUNNEL_HOME}/connectors/plugins-mapping.properties(仓库根目录下即 plugin-mapping.properties)中查到。也可以从 Apache Maven Repository 手动下载连接器 JAR 放入connectors/目录,效果相同。

步骤 3:添加作业配置文件定义作业

编辑config/v2.batch.config.template(仓库中的模板见 v2.batch.config.template),它决定了 SeaTunnel 启动后数据输入、处理和输出的方式及逻辑。下面是与示例链路完全一致的配置:

env { parallelism = 1 job.mode = "BATCH" } source { FakeSource { plugin_output = "fake" row.num = 16 schema = { fields { name = "string" age = "int" } } } } transform { FieldMapper { plugin_input = "fake" plugin_output = "fake1" field_mapper = { age = age name = new_name } } } sink { Console { plugin_input = "fake1" } }

对这份配置的关键参数做一点源码级说明:

  • env 区块parallelism控制作业并行度;job.mode"BATCH""STREAMING"。仓库自带的模板还展示了checkpoint.interval = 10000(毫秒)的写法,可作为流式作业参考。
  • FakeSource.row.num:每个并行度生成的数据条数。在 FakeSourceOptions.java 中可以看到row.num的默认值是 5,split.num(每个并行度生成的 split 数)默认 1,split.read-interval(两次 split 读取间隔,毫秒)默认 1。此外还支持rows(按行列表指定输出内容)、string.templateint.template等模板类参数,以及int.min/int.max等数值范围参数,方便按需构造更贴近真实场景的测试数据。
  • plugin_output/plugin_input:SeaTunnel 通过表名将上游输出与下游输入串联起来。FakeSource输出名为fakeFieldMapperfake读取、输出到fake1Console再从fake1读取,形成一条完整的数据流。
  • FieldMapper.field_mapper:字段映射与重命名规则,这里将name重命名为new_nameage保持不变。
  • Consolesink:打印到日志。其可配置项定义在 ConsoleSinkOptions.java,其中log.print.data默认true(是否打印数据),log.print.delay.ms默认 0(每条数据打印间隔毫秒)。

关于配置的更多信息可查看 配置的基本概念。

步骤 4:运行 SeaTunnel 应用程序

通过以下命令启动应用:

:::tip 从 2.3.1 版本开始,seatunnel.sh中的-e参数已被废弃,请改用-m参数。-m local对应的引擎模式定义在 MasterType.java(LOCAL("local"))。 :::

cd "apache-seatunnel-${version}" ./bin/seatunnel.sh --config ./config/v2.batch.config.template -m local

查看输出:运行命令后,控制台中打印的内容即是命令运行成功或失败的标志。SeaTunnel 控制台会打印类似下面的日志:

2022-12-19 11:01:45,417 INFO org.apache.seatunnel.connectors.seatunnel.console.sink.ConsoleSinkWriter - output rowType: new_name<STRING>, age<INT> 2022-12-19 11:01:46,489 INFO org.apache.seatunnel.connectors.seatunnel.console.sink.ConsoleSinkWriter - subtaskIndex=0 rowIndex=1: SeaTunnelRow#tableId=-1 SeaTunnelRow#kind=INSERT: CpiOd, 8520946 2022-12-19 11:01:46,490 INFO org.apache.seatunnel.connectors.seatunnel.console.sink.ConsoleSinkWriter - subtaskIndex=0 rowIndex=2: SeaTunnelRow#tableId=-1 SeaTunnelRow#kind=INSERT: eQqTs, 1256802974 2022-12-19 11:01:46,490 INFO org.apache.seatunnel.connectors.seatunnel.console.sink.ConsoleSinkWriter - subtaskIndex=0 rowIndex=3: SeaTunnelRow#tableId=-1 SeaTunnelRow#kind=INSERT: UsRgO, 2053193072 2022-12-19 11:01:46,490 INFO org.apache.seatunnel.connectors.seatunnel.console.sink.ConsoleSinkWriter - subtaskIndex=0 rowIndex=4: SeaTunnelRow#tableId=-1 SeaTunnelRow#kind=INSERT: jDQJj, 1993016602 2022-12-19 11:01:46,490 INFO org.apache.seatunnel.connectors.seatunnel.console.sink.ConsoleSinkWriter - subtaskIndex=0 rowIndex=5: SeaTunnelRow#tableId=-1 SeaTunnelRow#kind=INSERT: rqdKp, 1392682764 2022-12-19 11:01:46,490 INFO org.apache.seatunnel.connectors.seatunnel.console.sink.ConsoleSinkWriter - subtaskIndex=0 rowIndex=6: SeaTunnelRow#tableId=-1 SeaTunnelRow#kind=INSERT: wCoWN, 986999925 2022-12-19 11:01:46,490 INFO org.apache.seatunnel.connectors.seatunnel.console.sink.ConsoleSinkWriter - subtaskIndex=0 rowIndex=7: SeaTunnelRow#tableId=-1 SeaTunnelRow#kind=INSERT: qomTU, 72775247 2022-12-19 11:01:46,490 INFO org.apache.seatunnel.connectors.seatunnel.console.sink.ConsoleSinkWriter - subtaskIndex=0 rowIndex=8: SeaTunnelRow#tableId=-1 SeaTunnelRow#kind=INSERT: jcqXR, 1074529204 2022-12-19 11:01:46,490 INFO org.apache.seatunnel.connectors.seatunnel.console.sink.ConsoleSinkWriter - subtaskIndex=0 rowIndex=9: SeaTunnelRow#tableId=-1 SeaTunnelRow#kind=INSERT: AkWIO, 1961723427 2022-12-19 11:01:46,490 INFO org.apache.seatunnel.connectors.seatunnel.console.sink.ConsoleSinkWriter - subtaskIndex=0 rowIndex=10: SeaTunnelRow#tableId=-1 SeaTunnelRow#kind=INSERT: hBoib, 929089763 2022-12-19 11:01:46,490 INFO org.apache.seatunnel.connectors.seatunnel.console.sink.ConsoleSinkWriter - subtaskIndex=0 rowIndex=11: SeaTunnelRow#tableId=-1 SeaTunnelRow#kind=INSERT: GSvzm, 827085798 2022-12-19 11:01:46,491 INFO org.apache.seatunnel.connectors.seatunnel.console.sink.ConsoleSinkWriter - subtaskIndex=0 rowIndex=12: SeaTunnelRow#tableId=-1 SeaTunnelRow#kind=INSERT: NNAYI, 94307133 2022-12-19 11:01:46,491 INFO org.apache.seatunnel.connectors.seatunnel.console.sink.ConsoleSinkWriter - subtaskIndex=0 rowIndex=13: SeaTunnelRow#tableId=-1 SeaTunnelRow#kind=INSERT: EexFl, 1823689599 2022-12-19 11:01:46,491 INFO org.apache.seatunnel.connectors.seatunnel.console.sink.ConsoleSinkWriter - subtaskIndex=0 rowIndex=14: SeaTunnelRow#tableId=-1 SeaTunnelRow#kind=INSERT: CBXUb, 869582787 2022-12-19 11:01:46,491 INFO org.apache.seatunnel.connectors.seatunnel.console.sink.ConsoleSinkWriter - subtaskIndex=0 rowIndex=15: SeaTunnelRow#tableId=-1 SeaTunnelRow#kind=INSERT: Wbxtm, 1469371353 2022-12-19 11:01:46,491 INFO org.apache.seatunnel.connectors.seatunnel.console.sink.ConsoleSinkWriter - subtaskIndex=0 rowIndex=16: SeaTunnelRow#tableId=-1 SeaTunnelRow#kind=INSERT: mIJDt, 995616438

验证要点:

  • 首行output rowType: new_name<STRING>, age<INT>说明FieldMapper的字段重命名已生效;
  • 后续 16 行ConsoleSinkWriter输出对应row.num = 16生成的 16 条数据(字段已被重命名为new_nameage);
  • 批任务在写完全部数据后正常退出,进程结束。

Local 模式的行为与运维注意点

Local 模式下每个任务都会启动一个独立的进程,任务运行完成后进程退出。该模式有以下限制(参见 local-mode-deployment.md):

  1. 不支持任务的暂停、恢复;
  2. 不支持获取任务列表查看;
  3. 不支持通过命令取消作业,只能通过 Kill 进程的方式终止任务。

但每个任务由单独进程控制,不会出现任务之间相互影响的情况,适合对任务稳定性有强烈要求的场景。相关运维细节:

  • 运行日志输出到提交作业进程的标准输出;
  • 如需调整 JVM 参数,可修改$SEATUNNEL_HOME/config/jvm_client_options(该文件中的参数会应用到所有通过seatunnel.sh提交的作业,包括 Local 与集群模式),也可以在启动时追加,例如:./bin/seatunnel.sh --config ./config/v2.batch.config.template -m local -DJvmOption="-Xms2G -Xmx2G"

扩展示例:从 MySQL 到 Doris 批处理模式

跑通最小链路后,把 Source 和 Sink 换成真实连接器即可处理真实业务数据。下面以经典的 MySQL 到 Doris 批同步为例。

步骤 1:下载连接器

${SEATUNNEL_HOME}/config/plugin_config中加入连接器名称,然后执行安装命令(也可以从 Apache Maven Repository 手动下载连接器 JAR 放入connectors/目录),最后确认connector-jdbcconnector-doris都在${SEATUNNEL_HOME}/connectors/目录下。

# 配置连接器名称 --seatunnel-connectors-- connector-jdbc connector-doris --end--
# 安装连接器 sh bin/install-plugin.sh

步骤 2:放入 MySQL 驱动

下载 MySQL JDBC 驱动 JAR(mysql-connector-java),并放置在${SEATUNNEL_HOME}/lib/目录下,Jdbcsource 才能加载驱动建立连接。

步骤 3:添加作业配置文件定义作业

cd seatunnel/job/ vim st.conf
env { parallelism = 2 job.mode = "BATCH" } source { Jdbc { url = "jdbc:mysql://localhost:3306/test" driver = "com.mysql.cj.jdbc.Driver" connection_check_timeout_sec = 100 user = "user" password = "pwd" table_path = "test.table_name" query = "select * from test.table_name" } } sink { Doris { fenodes = "doris_ip:8030" username = "user" password = "pwd" database = "test_db" table = "table_name" sink.enable-2pc = "true" sink.label-prefix = "test-cdc" doris.config = { format = "json" read_json_by_line="true" } } }

参数说明:

  • Jdbc sourceurl为 MySQL 连接串,driver为驱动类名(8.x 驱动为com.mysql.cj.jdbc.Driver),user/password为账号口令,query为取数 SQL,table_path指定源表路径;
  • Doris sinkfenodes为 Doris FE 地址(ip:port),database/table为目标库表,sink.enable-2pc = "true"开启两阶段提交以保证写入一致性,sink.label-prefix设置事务标签前缀,doris.configformat = "json"read_json_by_line = "true"指定 JSON 按行流式写入。

关于配置的更多信息可查看 配置的基本概念。

步骤 4:运行 SeaTunnel 应用程序

cd seatunnel/ ./bin/seatunnel.sh --config ./job/st.conf -m local

查看输出:运行结束后,SeaTunnel 控制台会打印作业统计信息,作为成功或失败的标志:

*********************************************** Job Statistic Information *********************************************** Start Time : 2024-08-13 10:21:49 End Time : 2024-08-13 10:21:53 Total Time(s) : 4 Total Read Count : 1000 Total Write Count : 1000 Total Failed Count : 0 ***********************************************

Total Read CountTotal Write Count均为 1000 且Total Failed Count为 0,说明 1000 条数据已完整同步到 Doris。如需进一步优化作业,请参照对应连接器的使用文档调整参数。

第二部分:集群部署

如果已完成单机验证,希望在多节点环境中运行 SeaTunnel Engine,请继续阅读 SeaTunnel Engine(Zeta) 安装部署。集群部署文档集中说明了以下内容:

  • 不同部署模式的适用场景,包括 Local 模式、混合集群模式和分离集群模式;
  • 混合集群模式与分离集群模式的部署步骤;
  • 选择部署模式时的建议。

建议:

  • 如果只是想在一台机器上快速验证配置和任务链路,使用本文中的 Local 模式即可;
  • 如果需要多节点运行、资源隔离或更贴近测试和生产环境的部署方式,请进入集群部署文档继续操作。

下一步

  • 如果想先建立整体路径感,可以返回阅读 快速入门总览;
  • 当准备把示例 Source 和 Sink 替换成真实连接器时,建议继续阅读 作业配置指南;
  • 想直接看端到端核对的教程,可以先看 MySQL CDC 到 Kafka,再按链路形态选择 MySQL CDC 到 Doris、JDBC 到 S3、Kafka 到 Iceberg、Http 到 JDBC、File 到 StarRocks 和 多表 CDC;
  • 开始编写自己的配置文件,选择想要的连接器,并根据连接器文档配置参数;
  • 如果要部署多节点 SeaTunnel Engine 集群,请继续阅读 SeaTunnel Engine(Zeta) 安装部署;
  • 如果想进一步了解 SeaTunnel Engine,请参阅 SeaTunnel 引擎。

【免费下载链接】seatunnelSeaTunnel is a multimodal, high-performance, distributed, massive data integration tool.项目地址: https://gitcode.com/GitHub_Trending/se/seatunnel

创作声明:本文部分内容由AI辅助生成(AIGC),仅供参考

需要专业的网站建设服务?

联系我们获取免费的网站建设咨询和方案报价,让我们帮助您实现业务目标

立即咨询