
Apache SeaTunnel Engine 本地快速上手从单机验证 FakeSource 管线到 MySQL 批量入 Doris【免费下载链接】seatunnelSeaTunnel is a multimodal, high-performance, distributed, massive data integration tool.项目地址: https://gitcode.com/GitHub_Trending/se/seatunnel本文是 Apache SeaTunnel 内置引擎SeaTunnel Engine / Zeta的本地快速启动指南覆盖两条路径先用-m local在单机上验证安装、插件与作业配置再在需要多节点执行时平滑迁移到集群部署。读完本文你可以独立完成从部署 SeaTunnel、安装连接器插件、编写source-transform-sink三段式 HOCON 作业配置、运行批任务并解读输出日志到扩展为 MySQL → Doris 真实批量同步的完整过程。先理解两条路径SeaTunnel Engine 既可以用于单机快速试用也可以组成多节点集群运行。本文页面围绕这两条路径组织请根据你的目标选择路径适用场景下一步单机快速启动在单台机器上验证配置、连接器或作业管线继续阅读本文「单机快速启动本地模式」一节集群部署在测试、预发或类生产环境中跨多节点运行 SeaTunnel Engine前往 SeaTunnel Engine(Zeta) 部署指南建议当你想在单机上验证配置与作业管线时使用本文的本地模式当你需要多节点执行、资源隔离或更接近预发/生产的环境时再使用集群部署指南。Part 1单机快速启动本地模式本路径用于在单台机器上验证安装、连接器与作业配置。下面所有命令都以-m local方式启动 SeaTunnel Engine。开始前的准备工作如果你是第一次接触 SeaTunnel建议按顺序先阅读以下文档Getting Started 总览Deployment 部署说明Job 配置指南本文的示例管线使用了FakeSource模拟数据源、FieldMapper字段重命名转换和Console控制台输出 Sink三个插件。运行示例前请确保所需插件已安装。在${SEATUNNEL_HOME}/config/plugin_config中声明需要的连接器插件名然后执行安装脚本--seatunnel-connectors-- connector-fake connector-console --end--sh bin/install-plugin.sh在 Windows 上使用批处理脚本bin\install-plugin.cmd补充说明仓库根目录的 config/plugin_config 是当前发行版中插件声明文件的真实样例其中区块标记为--connectors-v2--自 2.3.x 起而非早期文档中的--seatunnel-connectors--你只需按所用发行版文档对应的标记格式填写插件名即可。所有受支持的连接器及其在plugin_config中的对应名称可在${SEATUNNEL_HOME}/connectors/plugins-mapping.properties中查到。第 1 步部署 SeaTunnel 与连接器开始之前请确保已按 Deployment 部署说明 下载并部署好 SeaTunnel 发行包需要安装 Java 8 或 11理论上高于 Java 8 的版本也可行并设置JAVA_HOME下载apache-seatunnel-version-bin.tar.gz二进制包并解压从 2.2.0-beta 起二进制包默认不再携带连接器依赖首次使用必须先运行sh bin/install-plugin.sh或 Windows 下的bin\install-plugin.cmd安装连接器也可以从 Apache Maven 仓库手动下载连接器 JAR 放入connectors/目录2.3.5 之前的版本放在connectors/seatunnel目录。如果你已经安装了全部连接器可以保留现有环境如果希望第一次运行尽量精简上文声明的connector-fake与connector-console两个插件足以支撑本页示例作业。第 2 步编写作业配置文件定义作业编辑config/v2.batch.config.template它决定了 SeaTunnel 启动后数据输入、处理与输出的方式和逻辑。仓库根目录的 config/v2.batch.config.template 是随发行版提供的同款模板文件。以下是与上述示例应用一致的配置文件内容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 1指定作业并行度为 1job.mode BATCH声明批模式。sourceFakeSource模拟生成数据通过plugin_output fake将输出表命名为fakerow.num 16表示生成 16 条数据schema.fields声明了namestring 类型与ageint 类型两个字段。transformFieldMapper读取输入表fakeplugin_input按field_mapper映射规则处理age age表示字段保持原名name new_name表示把name重命名为new_name处理结果输出到表fake1plugin_output。sinkConsole消费输入表fake1plugin_input把每一行数据打印到控制台。关于配置体系的更多信息可参考 配置基础概念。结合源码可以进一步理解各配置项的底层语义row.num的含义是每个并行度生成的数据条数定义于 FakeSourceOptions.java默认值为 5。示例中parallelism 1、row.num 16因此共输出 16 行数据。FakeSource 是“有界还是无界”取决于作业模式在 FakeSource.java 中可以看到BATCH模式返回Boundedness.BOUNDED有界流模式则返回UNBOUNDED无界。Console Sink 支持log.print.data是否打印数据默认true与log.print.delay.ms每条数据打印间隔毫秒数默认 0定义于 ConsoleSinkOptions.java。FakeSource 还支持split.num、split.read-interval、string.length、rows逐行指定模板数据、各类型的*.min/*.max/*.template/*.fake.mode等大量参数可用于构造更贴近真实业务的测试数据详见 FakeSourceOptions.java。FieldMapper转换插件位于seatunnel-transforms-v2模块的 fieldmapper 包 下除字段重命名外还支持按映射删除字段等能力。第 3 步运行 SeaTunnel 应用使用以下命令启动应用:::tip自 2.3.1 版本起seatunnel.sh中的-e参数已被弃用请改用-m。:::cd apache-seatunnel-${version} ./bin/seatunnel.sh --config ./config/v2.batch.config.template -m local在 Windows 上从 SeaTunnel 目录运行对应的批处理入口cd apache-seatunnel-3.0.0 bin\seatunnel.cmd --config config\v2.batch.config.template -m local观察输出命令运行后控制台会打印运行日志这是判断命令是否执行成功的重要信号。从源码层面看-m/--master参数由 ClientCommandArgs.java 定义支持local与cluster两种取值默认值为cluster其中-e与--deploy-mode已在 2.3.1 起标记为 deprecated。此外本地模式下 SeaTunnel 引擎会以MASTER_AND_WORKER角色嵌入进程运行参见 ClientExecuteCommand.java 中关于 local mode 的处理逻辑并在默认情况下暴露 REST/UI 端点便于你查看作业状态。SeaTunnel 控制台会打印类似如下的日志2022-12-19 11:01:45,417 INFO org.apache.seatunnel.connectors.seatunnel.console.sink.ConsoleSinkWriter - output rowType: new_nameSTRING, ageINT 2022-12-19 11:01:46,489 INFO org.apache.seatunnel.connectors.seatunnel.console.sink.ConsoleSinkWriter - subtaskIndex0 rowIndex1: SeaTunnelRow#tableId-1 SeaTunnelRow#kindINSERT: CpiOd, 8520946 2022-12-19 11:01:46,490 INFO org.apache.seatunnel.connectors.seatunnel.console.sink.ConsoleSinkWriter - subtaskIndex0 rowIndex2: SeaTunnelRow#tableId-1 SeaTunnelRow#kindINSERT: eQqTs, 1256802974 2022-12-19 11:01:46,490 INFO org.apache.seatunnel.connectors.seatunnel.console.sink.ConsoleSinkWriter - subtaskIndex0 rowIndex3: SeaTunnelRow#tableId-1 SeaTunnelRow#kindINSERT: UsRgO, 2053193072 2022-12-19 11:01:46,490 INFO org.apache.seatunnel.connectors.seatunnel.console.sink.ConsoleSinkWriter - subtaskIndex0 rowIndex4: SeaTunnelRow#tableId-1 SeaTunnelRow#kindINSERT: jDQJj, 1993016602 2022-12-19 11:01:46,490 INFO org.apache.seatunnel.connectors.seatunnel.console.sink.ConsoleSinkWriter - subtaskIndex0 rowIndex5: SeaTunnelRow#tableId-1 SeaTunnelRow#kindINSERT: rqdKp, 1392682764 2022-12-19 11:01:46,490 INFO org.apache.seatunnel.connectors.seatunnel.console.sink.ConsoleSinkWriter - subtaskIndex0 rowIndex6: SeaTunnelRow#tableId-1 SeaTunnelRow#kindINSERT: wCoWN, 986999925 2022-12-19 11:01:46,490 INFO org.apache.seatunnel.connectors.seatunnel.console.sink.ConsoleSinkWriter - subtaskIndex0 rowIndex7: SeaTunnelRow#tableId-1 SeaTunnelRow#kindINSERT: qomTU, 72775247 2022-12-19 11:01:46,490 INFO org.apache.seatunnel.connectors.seatunnel.console.sink.ConsoleSinkWriter - subtaskIndex0 rowIndex8: SeaTunnelRow#tableId-1 SeaTunnelRow#kindINSERT: jcqXR, 1074529204 2022-12-19 11:01:46,490 INFO org.apache.seatunnel.connectors.seatunnel.console.sink.ConsoleSinkWriter - subtaskIndex0 rowIndex9: SeaTunnelRow#tableId-1 SeaTunnelRow#kindINSERT: AkWIO, 1961723427 2022-12-19 11:01:46,490 INFO org.apache.seatunnel.connectors.seatunnel.console.sink.ConsoleSinkWriter - subtaskIndex0 rowIndex10: SeaTunnelRow#tableId-1 SeaTunnelRow#kindINSERT: hBoib, 929089763 2022-12-19 11:01:46,490 INFO org.apache.seatunnel.connectors.seatunnel.console.sink.ConsoleSinkWriter - subtaskIndex0 rowIndex11: SeaTunnelRow#tableId-1 SeaTunnelRow#kindINSERT: GSvzm, 827085798 2022-12-19 11:01:46,490 INFO org.apache.seatunnel.connectors.seatunnel.console.sink.ConsoleSinkWriter - subtaskIndex0 rowIndex12: SeaTunnelRow#tableId-1 SeaTunnelRow#kindINSERT: NNAYI, 94307133 2022-12-19 11:01:46,490 INFO org.apache.seatunnel.connectors.seatunnel.console.sink.ConsoleSinkWriter - subtaskIndex0 rowIndex13: SeaTunnelRow#tableId-1 SeaTunnelRow#kindINSERT: EexFl, 1823689599 2022-12-19 11:01:46,490 INFO org.apache.seatunnel.connectors.seatunnel.console.sink.ConsoleSinkWriter - subtaskIndex0 rowIndex14: SeaTunnelRow#tableId-1 SeaTunnelRow#kindINSERT: CBXUb, 869582787 2022-12-19 11:01:46,490 INFO org.apache.seatunnel.connectors.seatunnel.console.sink.ConsoleSinkWriter - subtaskIndex0 rowIndex15: SeaTunnelRow#tableId-1 SeaTunnelRow#kindINSERT: Wbxtm, 1469371353 2022-12-19 11:01:46,491 INFO org.apache.seatunnel.connectors.seatunnel.console.sink.ConsoleSinkWriter - subtaskIndex0 rowIndex16: SeaTunnelRow#tableId-1 SeaTunnelRow#kindINSERT: mIJDt, 995616438解读这些日志可以发现两个关键点首行日志打印了经过转换后的输出行类型output rowType: new_nameSTRING, ageINT——name字段已被FieldMapper重命名为new_name证明transform环节生效后续每行是ConsoleSinkWriter源码逐条输出的SeaTunnelRow每条数据包含INSERT行类型与两个字段值rowIndex从 1 到 16与row.num 16完全对应。扩展示例MySQL 到 Doris 的批模式同步本地模式验证通过后把示例中的模拟源与打印 Sink 替换为真实连接器即可投入实际同步。下面演示最常见的「MySQL 批量读入、Doris 批量写出」场景。第 1 步下载连接器首先将连接器名称添加到${SEATUNNEL_HOME}/config/plugin_config文件中然后执行命令安装连接器当然你也可以从 Apache Maven 仓库手动下载连接器放入connectors/目录。最后确保connector-jdbc和connector-doris两个连接器都位于${SEATUNNEL_HOME}/connectors/目录下。# 配置连接器名称。 --seatunnel-connectors-- connector-jdbc connector-doris --end--# 安装连接器。 sh bin/install-plugin.sh第 2 步放置 MySQL 驱动你需要下载 MySQL 的 JDBC 驱动 JAR 包并把它放到${SEATUNNEL_HOME}/lib/目录下以便 JDBC 连接器能够加载 MySQL 驱动。第 3 步添加作业配置文件定义作业cd seatunnel/job/ vim st.confenv { 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_linetrue } } }配置要点env.parallelism 2批任务以并行度 2 执行相比示例作业提升吞吐。JdbcSource通过url指定 MySQL 连接串、driver指定驱动类com.mysql.cj.jdbc.Driver、connection_check_timeout_sec设置连接检查超时秒、user/password设置账号密码、table_path指定库表、query指定读取 SQL。DorisSinkfenodes为 Doris FE 地址ip:portusername/password为账号密码database/table为目标库表sink.enable-2pc true开启两阶段提交保证精确一次语义sink.label-prefix设置导入 label 前缀doris.config内通过format 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 ***********************************************这份「Job Statistic Information」是批作业最直接的验收证据Total Read Count与Total Write Count均为 1000Total Failed Count为 0说明 MySQL 中的 1000 行数据在 4 秒内全部写入 Doris且无失败记录。:::tip如果希望优化作业可参考连接器文档Source-MySQL 与 Sink-Doris。:::Part 2集群部署如果你已经在本地验证了作业并希望跨多个节点运行 SeaTunnel Engine请继续阅读 SeaTunnel Engine(Zeta) 部署指南。该部署指南覆盖以下内容本地模式、混合集群模式与分离集群模式的部署场景混合集群模式与分离集群模式的具体部署步骤如何选择合适的部署模式。建议当你想在单机上验证配置与作业管线时使用本页的本地模式当你需要多节点执行、资源隔离或更接近预发/生产的环境时再使用部署指南。延伸阅读想要一条贯穿部署、快速启动与配置的引导式阅读路径从 Getting Started 总览 开始。准备把示例中的 Source 与 Sink 替换为真实连接器时参考 Job 配置指南。希望继续走一遍经过验证的「源到目标」演练可以接着看 MySQL CDC 到 Kafka其他管线形态见 MySQL CDC 到 Doris、JDBC 到 S3、Kafka 到 Iceberg、Http 到 JDBC、File 到 StarRocks 或 多表 CDC。开始编写你自己的配置文件选择想用的 连接器并根据连接器文档配置参数。想要部署多节点 SeaTunnel Engine 集群继续阅读 SeaTunnel Engine(Zeta) 部署指南。想进一步了解 SeaTunnel Engine 本身可阅读 SeaTunnel Engine(Zeta) 介绍。【免费下载链接】seatunnelSeaTunnel is a multimodal, high-performance, distributed, massive data integration tool.项目地址: https://gitcode.com/GitHub_Trending/se/seatunnel创作声明:本文部分内容由AI辅助生成(AIGC),仅供参考