
简介面向数据工程师与开发人员的Apache SeaTunnel可运行配置案例聚焦MySQL到HDFS、Hive到MySQL两条数据迁移链路提供完整配置示例、前置操作指引与常见问题排查思路适合有一定大数据基础、需要快速落地数据集成任务的读者。压缩包共3个文件分别为HTML配置说明页、可运行的云端环境配置inscode以及Git忽略规则配置文件整体仅8KB结构精简清晰。目前已有27人学习。案例从创建目录、放置MySQL连接器jar包、启动Hadoop等准备工作讲起再逐步展开MySQL到HDFS的两种配置写法以及Hive到MySQL转换时的数据准备、服务启动与执行命令要点同时点出各环节易错点。下载后即可对照源码修改参数在真实环境中验证数据迁移效果显著缩短工具学习与配置调优周期。1. SeaTunnel 配置案例先搞清楚它到底解决什么问题如果你手头有十几个数据源每天要同步到数仓或业务库还时不时冒出一个新需求——把 MySQL 的数据同步到 Doris把 Kafka 的消息落到 Hive或者把某个 API 接口的数据攒进 ClickHouse那么你大概率已经受够了写一堆定制同步脚本的日子。SeaTunnel 就是干这个的一套配置搞定数据接入、转换、写出跑在 Flink 或 Spark 上也能本地直跑关键是它不需要你为每个数据源写一套连接器代码。这篇文章的价值在于不给你画架构大饼直接用可运行的配置源码把整个链路串起来。你读完能照着我给的配置改几个参数就能在你自己的环境里跑通 MySQL → 控制台的同步再扩展到 Kafka → Doris、JDBC 多表同步。我会把每个配置项的含义、为什么这么设、以及最容易翻车的地方都交代清楚。适合刚接触 SeaTunnel、被官方文档里一堆连接器参数搞得头大的人也适合已经在用但想在配置层面优化一把的开发者。我最早用 SeaTunnel 就是被多数据源同步逼的——每接一个源就要写一个 Flink 程序后来发现 SeaTunnel 把 90% 的工作压缩成了改 YAML剩下的 10% 才是真正需要动代码的地方。所以这篇文章的核心就一句话配置即代码跑通即交付。2. SeaTunnel 配置体系从 Job 定义到连接器选型2.1 一个最小可运行配置的完整拆解SeaTunnel 的配置核心是一个 Job 定义描述数据从哪来、经过什么处理、写到哪去。整个配置文件通常是hocon格式也有 JSON 格式看起来像类型安全的配置文件但里面藏着几个必须理解的关键块。先看一个最基础、保证能直接跑起来的配置。把 SeaTunnel 下载解压后在config目录新建一个first_job.confenv { parallelism 2 job.mode BATCH } source { FakeSource { result_table_name fake_tbl row.num 3 schema { fields { id int name string } } } } transform { # 这里先不配数据原样传递 } sink { Console { source_table_name fake_tbl } }这是 SeaTunnel 世界里的 Hello WorldFakeSource自己生成 3 条测试数据Consolesink 把它们打印到控制台。逻辑说明env块定义运行参数最常调的是parallelism决定并行度job.mode取值BATCH或STREAMING决定是跑一次还是持续跑。source、transform、sink是三层结构层内可以配多个连接器用result_table_name和source_table_name做数据流转的“接线”。schema里定义了字段名和类型FakeSource 就按这个 schema 造数。参数说明row.num控制造几条数据parallelism设太大会造成小任务资源浪费本地测试建议从 2 开始。这个配置跑通后你就明白 SeaTunnel 的核心抽象了Source 产数据Transform 改数据Sink 收数据表名就是数据管道里的“临时表”。后面所有复杂配置都是在这个骨架上加东西。2.2 Source 连接器怎么选Jdbc、Kafka、FakeSource 的取舍连接器选型直接决定你的同步任务稳不稳。把最常用的三个拿出来对比你就知道什么时候该用谁。连接器适用场景关键参数常见坑Jdbc Source关系库离线同步、批量抽取url、user、password、query驱动要单独放分页参数要调Kafka Source实时流接入、日志/事件流bootstrap.servers、topic、consumer.groupIdoffset 重置策略、序列化格式FakeSource测试、压测、链路验证row.num、schema没有纯粹造数选型理由如果你的任务是一次性把 MySQL 某张表抽到数仓用 Jdbc Source 配合 SQL 查询就行没必要引入 Kafka 增加链路复杂度。如果数据源已经在 Kafka 里躺着那就别再用 JDBC 去查——直接用 Kafka Source 消费省掉中间落库环节。FakeSource 则适合你在不依赖真实环境的情况下验证下游 Sink 的配置是否正确。一个常见误用场景有人为了“统一架构”离线同步也把 MySQL 数据先发到 Kafka 再消费结果多了一层运维成本和数据延迟。我的习惯是离线批量用 Jdbc、实时流式用 Kafka、测试验证用 Fake。这不是什么金科玉律但能少踩很多坑。2.3 Sink 端连接器配置从 Console 到 Doris、MySQL 的跨越Sink 端是整个配置里最容易出问题的地方因为每个目标系统的写入协议、批次大小、事务机制都不一样。SeaTunnel 的 Sink 连接器大致分两类一个是极简的Console用于验证一类是真实写入型比如Jdbc、Doris、ClickHouse、Hive。以最常见的 MySQL Sink 为例sink { Jdbc { url jdbc:mysql://172.18.1.102:3306/demo?useUnicodetruecharacterEncodingUTF-8 driver com.mysql.cj.jdbc.Driver user root password root query INSERT INTO target_table (id, name) VALUES (?, ?) source_table_name fake_tbl } }逻辑说明query写的是预编译 SQL?占位符按字段顺序绑定上游数据。driver必须和你的 JDBC 驱动匹配MySQL 8 用com.mysql.cj.jdbc.DriverMySQL 5.x 用com.mysql.jdbc.Driver搞错直接报类找不到。参数说明url里的字符集参数建议显式加上否则中文写进去变成乱码user、password不说了注意不要跟其他任务混用账号排查问题时会很痛苦。如果目标表字段多、写入量大query还可以写成INSERT INTO table (a,b,c) VALUES (?,?,?) ON DUPLICATE KEY UPDATE ...来做 upsert。从 Console 跨到真实 Sink你只需要做三件事把 Console 换成具体的 JDBC/Doris 连接器、填好连接信息、写好 SQL 或指定表映射。这个过程里最容易出问题的就是驱动版本和 SQL 方言下面专门讲避坑。3. 可运行配置案例一套完整的 MySQL 到 Doris 同步配置3.1 环境准备下载、解压、启动的最小步骤这一步是照着做就能跑通的前提。先去 SeaTunnel 官网下载对应的发布包注意选择与你的 Flink/Spark 环境匹配的版本或者直接用自带的本地引擎然后解压到指定目录。# 以 2.3.x 版本为例解压后目录结构大致如下 tar -zxvf apache-seatunnel-2.3.x-bin.tar.gz cd apache-seatunnel-2.3.x # 如果你用自带引擎跑不需要额外启动 Flink 集群 # 先检查插件目录里有没有你要用的连接器 ls plugins/ # 期望看到 seatunnel-plugin-connector-jdbc-xxx.jar 这类文件 # 如果没有对应连接器需要先安装 sh bin/install-plugin.sh 2.3.x逻辑说明SeaTunnel 把连接器做成独立插件plugins目录下每个 jar 对应一个连接器。运行任务时SeaTunnel 会根据配置里的连接器名比如Jdbc、Doris去加载对应的插件。如果你下载的包是精简版可能缺少某些连接器那就用install-plugin.sh按版本号拉取。参数说明install-plugin.sh后面的版本号要和当前 SeaTunnel 版本一致。安装过程会从中央仓库下载 jar如果网络受限可能需要配置镜像或者手动把 jar 拷进plugins目录。启动本地任务的方式很简单# 本地跑 batch 任务 sh bin/seatunnel.sh --config config/mysql_to_doris.conf -e local-e local指定本地执行模式意味着不依赖外部 Flink/Spark 集群最适合开发调试。等你确认配置没问题再换-e spark或-e flink提交到集群。3.2 MySQL Source 到 Doris Sink 的完整配置源码这是本文最核心的可运行配置案例。场景把 MySQL 里的orders表增量同步到 Doris 的ods_orders表字段做了简单的类型映射和过滤。env { parallelism 2 job.mode BATCH } source { Jdbc { url jdbc:mysql://172.18.1.102:3306/business_db?useUnicodetruecharacterEncodingUTF-8zeroDateTimeBehaviorconvertToNull driver com.mysql.cj.jdbc.Driver user sync_user password sync_password query SELECT id, order_no, user_id, amount, create_time FROM orders WHERE create_time 2024-01-01 00:00:00 partition_column id partition_num 4 result_table_name mysql_orders } } transform { # 字段名保持原样实际项目中可以在这里做数据清洗 # 比如把 amount 从分转成元amount / 100 } sink { Doris { fenodes 172.18.1.103:8030 username doris_user password doris_password database ods_db table ods_orders source_table_name mysql_orders field_mapper { id id order_no order_no user_id user_id amount amount create_time create_time } sink.label.prefix seatunnel_mysql_orders doris.config { format json read_json_by_line true time_zone Asia/Shanghai } } }逻辑说明Source 端用 Jdbc 连接器执行一条查询 SQLpartition_column和partition_num是让 SeaTunnel 按 id 把查询拆成 4 个分片并行读取——如果表很大不加这个参数就只能单线程查性能差距极大。Sink 端用 Doris 连接器fenodes指向 Doris 的 FE 节点field_mapper做字段映射sink.label.prefix是 Doris Stream Load 的标签前缀用于保证导入事务的幂等性。参数说明MySQL 的zeroDateTimeBehaviorconvertToNull必须加否则遇到0000-00-00 00:00:00这种时间值会直接报错partition_num不是越大越好建议和表数据量、并行度匹配经验值是 48。doris.config里的formatjson配合read_json_by_linetrue是 SeaTunnel 写入 Doris 时最稳定的组合用 CSV 容易在特殊字符上踩坑。3.3 运行与验证从提交命令到控制台输出检查配置写好后运行命令如下sh bin/seatunnel.sh --config config/mysql_to_doris.conf -e local执行后正常会看到类似这样的日志2024-xx-xx 12:00:01 INFO SeatunnelJob - starting job: mysql_to_doris 2024-xx-xx 12:00:05 INFO SeaTunnel - source Jdbc reader: 10000 rows 2024-xx-xx 12:00:08 INFO SeaTunnel - sink Doris writer: 10000 rows 2024-xx-xx 12:00:08 INFO SeatunnelJob - job finished, status: SUCCEEDED验证数据是否写对去 Doris 里查一下SELECT COUNT(*) FROM ods_db.ods_orders; SELECT * FROM ods_db.ods_orders ORDER BY id DESC LIMIT 5;逻辑说明source Jdbc reader: 10000 rows表示 Source 读出来 10000 行sink Doris writer: 10000 rows表示 Sink 写进 Doris 也是 10000 行两边对上就说明配置基本没问题。参数说明如果 Sink 侧行数少于 Source优先查 Doris 的标签是否重复导致去重或者field_mapper里有没有字段类型不匹配被静默丢弃如果完全没写入查 Doris FE 的日志和 SeaTunnel 的错误堆栈重点看doris.config里的参数格式是否正确。跑通这个案例后你已经掌握了 SeaTunnel 最核心的用法选对 Source 连接器、配置查询、选对 Sink 连接器、映射字段、提交运行、检查行数。后面所有变更都是在这个流程上做参数级调整。4. 配置避坑指南驱动、并行度与连接串的翻车现场4.1 驱动缺失导致的 ClassNotFoundException 现象与解决现象配置没问题、服务也通但提交任务直接抛ClassNotFoundException: com.mysql.cj.jdbc.Driver。原因很简单——SeaTunnel 的发行包不内置每个数据库的 JDBC 驱动MySQL 驱动需要你自己放进plugins或lib目录。解决方法是# 先把驱动 jar 下载到本地注意版本要跟数据库服务端匹配 # mysql-connector-java-8.0.xx.jar 对应 MySQL 8.x cp mysql-connector-java-8.0.xx.jar $SEATUNNEL_HOME/lib/ # 然后重启任务这里要注意MySQL 驱动 jar 下载时别选成mysql-connector-j的新命名旧版的mysql-connector-java在 Maven 中央仓库也能找到两者类名一样但兼容性有细微差别。我的习惯是优先用和 SeaTunnel 版本发布时间接近的驱动版本避免因为驱动新特性导致连接参数不识别。4.2 并行度设置不合理导致任务 OOM 或源端被打满现象任务一跑SeaTunnel 内存飙到 90% 以上或者 MySQL 的 CPU 直接被拉满其他业务查询全变慢。原因parallelism和partition_num乘积过大比如parallelism5、partition_num10等于同时开了 50 个 JDBC 连接去查同一张表。解决思路env { parallelism 2 } source { Jdbc { partition_num 2 } }参数说明先按源库负载预估再按 SeaTunnel 所在机器的内存调整。简单粗暴的公式是partition_num不要超过源表数据量除以 100 万意思是每 100 万行一个分片比较合理。如果源表只有 10 万行partition_num1就够了。4.3 连接串缺参数导致中文乱码或时间解析失败现象同步到 Doris 后中文全部变成?或者??时间字段变成null。原因JDBC URL 里缺少字符集参数或者 Doris 侧的time_zone没设对。解决方式url jdbc:mysql://172.18.1.102:3306/business_db?useUnicodetruecharacterEncodingUTF-8zeroDateTimeBehaviorconvertToNulluseSSLfalseuseSSLfalse也很重要某些 MySQL 版本默认开启 SSL会导致连接变慢甚至握手失败。Doris 侧指定时区doris.config { time_zone Asia/Shanghai }4.4 使用可运行源码做二次开发时的包名冲突问题现象你从 GitHub 拉下来的 SeaTunnel 源码编译后运行报NoClassDefFoundError或者方法签名对不上。原因Maven 依赖里引入了跟 SeaTunnel 自带插件 jar 冲突的第三方包。解决方式开发调试时用mvn -pl seatunnel-examples -am package只构建需要的模块避免全量构建产生的 jar 冲突并且运行前检查lib目录里有没有多个版本的连接器 jar。这个坑在新手二次开发时特别常见本质是构建工具链和运行时隔离的问题。建议直接在 IDE 里跑调试模式这样能更快定位冲突来源而不是一次次打 jar 包试错。4.5 增量同步重复数据或漏数据的排查思路现象Count 校验时发现 Doris 里的行数比 MySQL 多或者少。原因增量同步的 WHERE 条件没有边界比如create_time 2024-01-01每次跑都会把 1 月 1 号的数据再插一遍漏数据则可能是 MySQL 事务隔离级别导致读取到未提交数据。解决对于重复Sink 侧用INSERT INTO ... ON DUPLICATE KEY UPDATE做幂等对于漏数据Source 侧查询语句加上主键或唯一时间戳的上限比如AND create_time NOW()。这是同步任务里最常见的两个坑本质是没把“增量”的逻辑做严谨。5. Transform 配置实操过滤、转换与多表数据拼接5.1 用 Filter 和 Replace 做数据清洗的配置片段SeaTunnel 的 Transform 层在字数上经常被忽略但实际项目里它是减少下游加工负担的关键。比如从 MySQL 同步过来的user_phone字段需要把空格和86去掉transform { Filter { source_table_name mysql_orders result_table_name orders_clean filter_expressions [ amount 0 ] } Replace { source_table_name orders_clean result_table_name orders_replaced fields [user_phone] replace { user_phone { pattern [ 86] replacement } } } }逻辑说明Filter相当于 SQL 的 WHERE只让满足filter_expressions的行往下走。Replace是对指定字段做正则替换pattern是目标模式replacement是替换成的字符串。这里把user_phone中的空格、加号、数字 86 全部去掉就得到了干净的手机号。参数说明filter_expressions支持比较运算符和逻辑与、或||组合。多个 Transform 之间用result_table_name和下一个的source_table_name串起来数据流像管道一样一节节往前传。注意Replace的正则语法走的是 Java Pattern 规则转义符要写双反斜杠比如匹配点号要写\\.。5.2 用 SQL Transform 完成行转列和聚合计算除了简单的 Replace、FilterSeaTunnel 还内置了一个Sql多表拼接和聚合的 Transform。它允许你用 SQL 直接处理上游表这对习惯写 SQL 的人非常友好transform { Sql { source_table_name orders_clean result_table_name orders_agg query SELECT user_id, COUNT(*) AS order_cnt, SUM(amount) AS total_amount FROM orders_clean GROUP BY user_id } }逻辑说明上面这个配置把orders_clean表按user_id聚合统计每个用户的订单数和总金额结果写入orders_agg表然后下游 Sink 可以消费这张聚合表。Sqltransform 底层本质上是在 SeaTunnel 引擎里启动了一个微型 SQL 执行器目前支持常用的聚合、JOIN、子查询但不支持窗口函数和复杂 UDF。参数说明query里的表名必须对应前面的result_table_name不能直接写 MySQL 里的物理表名。SQL transform 处理的数据集是批式的如果你是流式任务它会在每个微批次上执行这个 SQL因此聚合结果是按批次粒度计算的不是全量全局聚合。这一点在实时同步场景要特别留意别把实时数据当成离线表来聚合。5.3 多表拼接把两张 MySQL 表合成一个 Sink 结果集多表拼接也是高频需求比如主表orders和字典表order_status_dict需要 JOIN 后写入 Doris 宽表。用 SQL transform 可以这样实现source { Jdbc { result_table_name orders query SELECT id, status_code, amount FROM orders # ... 连接参数同上 } Jdbc { result_table_name status_dict query SELECT code, name FROM order_status_dict # ... 连接参数同上 } } transform { Sql { result_table_name orders_enriched query SELECT o.id, o.amount, d.name AS status_name FROM orders o LEFT JOIN status_dict d ON o.status_code d.code } }注意这里有两个 Jdbc Source 配置块用不同的result_table_name区分两路输入然后在Sqltransform 里直接用表名 JOIN。Sink 端指向orders_enriched即可。这个方案非常适合“主数据 维表补齐”的场景避免在 Doris 里做昂贵的 JOIN。但要注意如果维表很大百万级以上每次任务都全量加载维表代价很高。我的经验是维表控制在 10 万行以内否则想办法预处理成小文件或者走维表关联插件。6. 端到端验证与进阶从配置漂移到集群调优的工程化习惯当你把 SeaTunnel 的配置从“能跑”推向“稳定跑”真正拉开差距的往往不是某个连接器参数而是一套验证和治理习惯。我自己的交付标准是配置必须有版本管理、任务必须有水位线、数据必须能被重复校验。一个具体做法是把 SeaTunnel 目录下的所有.conf文件纳入 Git 仓库。每次改动配置提交信息里写清楚“改了哪个 Sink 的 batch 大小”“为什么调低并行度”。等生产环境出问题你能靠git diff快速定位是哪一次提交改了什么。# 在 config 目录初始化 Git 仓库 cd $SEATUNNEL_HOME/config git init git add *.conf git commit -m init seatunnel configs # 后续每次改配置都提交再说验证。我一般会做一个“影子表”验证策略在 Doris 里建一张ods_orders_shadow表结构和正式表一模一样先跑一次全量同步到影子表对比行数、关键字段的 SUM 和 COUNT都通过后再切换到正式表。这样能避免源端脏数据直接毒害生产表。进阶用户可能还会关注流式场景的 checkpoint 配置比如 Kafka Source 配合STREAMING模式时env里要显式配checkpoint.interval默认 10 秒。设置太短Sink 频繁提交导致吞吐下降设置太长故障恢复时数据回放窗口变大。我的取值习惯是 3060 秒在吞吐和延迟之间取平衡。最后是一个真实教训有次我把 MySQL 大表同步到 Doris初版配置没开partition_column任务跑了 2 小时没跑完。后来加上主键分片、并行度调到 4同样的数据量 20 分钟就结束了。从此我养成了一个习惯——任何超过 10 万行的源表第一件事就是问自己分片了吗并行度够不够这个习惯帮我省下的时间远超早期在集成上花的时间。希望这些配置案例和踩坑记录能帮你把 SeaTunnel 真正用起来少走几步弯路。本文还有配套的精品资源点击获取