简介:面向数据工程师与开发人员的Apache SeaTunnel可运行配置案例包,聚焦MySQL到HDFS、Hive到MySQL两类常见数据迁移场景,适合正在搭建数据集成流程或需要参考完整配置结构的初中级使用者。包体为zip压缩包,共3个文件,包含inscode在线运行工程配置、html说明页面与gitignore文件,整体仅8KB,轻量易用,解压后即可对照阅读。已有27人学习。案例包内不仅给出源端与目标端参数设置要点,还涵盖目录创建、连接器jar包放置、Hadoop/Hive服务启动及转换命令执行等前置与收尾细节,并配有可运行源码便于直接验证。通过阅读html说明和导入inscode工程,能快速理解SeaTunnel配置文件的写法与不同参数对数据迁移的实际影响,节省自行摸索时间。
1. 先把这套 SeaTunnel 配置案例跑通:从零到第一个同步任务
SeaTunnel 配置案例(可运行源码)解决的不是“SeaTunnel 怎么安装”,而是“配置怎么写才能一次跑通”。做过数据集成的人都有这种经历:照着文档抄、连接串看着也没问题,任务就是起不来。目录里少一个 jar、URL 少一个时区参数、sink 的 SQL 列顺序跟表结构错位,任何一个都能让同步任务在第一个 checkpoint 之前翻车。这套资源把 MySQL 全量、MySQL-CDC 增量、Kafka 消费、ClickHouse 落地四个高频场景的可运行配置和启动脚本打包在一起,新手可以从零跑到出数,熟手可以直接把模板改成自己的业务场景。
2. 从环境到跑通第一个 Job:目录、引擎与最小配置
2.1 先弄懂发行包三个目录,不然报错都看不懂
拿到 SeaTunnel 发行包解压后,我的习惯是先不看文档,先把bin/、lib/、connectors/(老版本叫plugins/)三个目录完整看一遍。bin/放启动脚本,lib/放引擎核心 jar 和 JDBC 驱动,connectors/下按插件名分目录存放各个连接器 jar。这三个目录决定了后面大多数报错的排查方向:报ClassNotFoundException去lib/查驱动,报Plugin not found去connectors/查连接器,报语法错误才回头查.conf文件。很多人一上来就改配置,改到半夜也没用,因为问题根本不在配置本身。
引擎选型上,这套案例默认走 Zeta 引擎,也就是 SeaTunnel 自带的分布式引擎,不需要额外部署 Flink 或 Spark。对单机测试、中小规模同步、想快速出数的场景,Zeta 是成本最低的选择:一个 Java 进程,一条命令启动,配置里写job.mode就行。如果环境里已经有一套 Flink 集群,并且希望所有数据任务统一由 Flink 调度,再考虑用 Flink 引擎的发行包。配置文件的写法差别不大,但启动命令、资源申请方式、日志入口完全不同。我的建议是:没有存量集群就别为了“看起来更分布式”去硬套 Flink,单机跑 Zeta 能解决的问题,不值得引入一套调度框架。
配置文件是 HOCON 风格,长得像带缩进的 JSON,但脾气比 JSON 大。缩进用空格,不要用 Tab,否则解析器很可能在某一行缩进上直接报错;字符串值大部分时候可以裸写,但值里出现:、,、#这类特殊字符时,最好用双引号包住。案例包里的.conf文件都遵循同一个结构:env段写并行度与任务模式,source段定义数据从哪来,transform段做可选加工,sink段定义数据写到哪。四个段里transform可以省略,但段与段的顺序不能乱,乱了的报错信息非常不直观,属于能卡住新手半小时的“黑匣子”。
2.2 最小可运行配置:MySQL 全量同步到 MySQL
案例包里排第一的是mysql_to_mysql.conf,最典型的入门任务:把一张表全量搬到另一张表。下面是精简后的核心配置,去掉了注释和多余空行,实际源码包里每个参数都带说明。
env { parallelism = 2 job.mode = "BATCH" } source { Jdbc { url = "jdbc:mysql://127.0.0.1:3306/demo_db?serverTimezone=Asia/Shanghai&useSSL=false&characterEncoding=UTF-8" user = "root" password = "123456" query = "SELECT id, name, created_at FROM demo_user WHERE id > 0" result_table_name = "src_user" } } sink { Jdbc { url = "jdbc:mysql://127.0.0.1:3306/demo_db?serverTimezone=Asia/Shanghai&useSSL=false&characterEncoding=UTF-8" user = "root" password = "123456" query = "INSERT INTO demo_user_copy(id, name, created_at) VALUES (?, ?, ?)" } }逻辑说明:source段里Jdbc插件连源库,query按过滤条件查出要同步的数据,result_table_name是这一路数据的输出别名。sink段里Jdbc插件连目标库,query是预编译写入语句,?的数量和顺序必须和 source 输出的列一一对应。这个任务我刻意不写transform段,数据原样搬过去;job.mode = "BATCH"表示跑完就退出,适合定时全量同步。如果把它改成STREAMING,任务会常驻,那是后面 CDC 案例的模式。
参数说明:serverTimezone=Asia/Shanghai解决 MySQL 8 驱动时区报错;characterEncoding=UTF-8避免中文在传输过程中变成问号;useSSL=false省掉本地调试时的 SSL 握手开销。parallelism = 2在 batch 任务里表示两个并行读源,实际取值要结合源库负载调整,这不是越大越好,后面第 5 章专门讲。WHERE id > 0是我刻意加的,防止有人把全表无过滤的 SQL 直接搬去生产,一跑就是一次全表扫描,源库连接池直接被打满。
2.3 启动命令与日志验证:怎么确认 Job 真的跑完了
配置写好后,启动命令只有三条:
cd /opt/apache-seatunnel # 驱动必须放 lib 目录,JDBC 源和 JDBC sink 都靠它才能加载 cp mysql-connector-j-8.0.33.jar lib/ # 启动 batch 同步任务 bin/seatunnel.sh --config config/mysql_to_mysql.conf说明:新版本 Zeta 引擎发行包默认本地模式运行,--config指定配置文件路径;老版本部分发行包要求加-e local或-m local,具体以bin/seatunnel.sh --help输出为准,不同小版本命令略有差异。驱动 jar 的版本要和 MySQL 服务端匹配,MySQL 8.x 用com.mysql.cj.jdbc.Driver,5.x 用com.mysql.jdbc.Driver,放错版本会直接报驱动类找不到。
跑起来之后不要只盯着“有没有报错”,要确认“到底跑了多少行”。任务正常结束时,控制台日志里会出现JOB_FINISHED字样,模块统计里能看到writeRecordCount,那个数字是本批次的写入行数。我一般会拿它和源表COUNT(*)对一遍,对不上就先查query里的WHERE条件有没有写错,而不是急着翻目标表数据。第一次跑不通时,不要反复改配置试运气,先看日志里第一个Exception的堆栈第一行。多数情况是驱动没加载,少数情况是 SQL 语法被目标库拒绝,这两种错的排查方向完全不同,混在一起只会越改越乱。
提示:跑批量任务前,先在测试库建一张空表,用
EXPLAIN确认源端 query 不会全表扫描,再上生产。这个习惯能帮你躲掉一半以上的源库故障。
3. 四个高频场景配置模板拆解:批量、增量、流式与多路写入
3.1 MySQL-CDC 增量同步:initial 模式同时拿快照和 binlog
全量同步能跑通之后,接下来需求一定是增量。MySQL-CDC 是 SeaTunnel 里最常用的增量方案,核心思路是:先做一次全量快照,再切换成 binlog 实时读取。案例包里的mysql_cdc_to_mysql.conf把这两件事合并到了一个任务里,配置关键点如下:
env { parallelism = 1 job.mode = "STREAMING" } source { MySQL-CDC { result_table_name = "cdc_user" host = "127.0.0.1" port = 3306 username = "root" password = "123456" database-name = "demo_db" table-names = ["demo_db.demo_user"] server-id = "5400-5404" startup.mode = "initial" } } sink { Jdbc { url = "jdbc:mysql://127.0.0.1:3306/demo_db?serverTimezone=Asia/Shanghai&useSSL=false" user = "root" password = "123456" query = "INSERT INTO demo_user_binlog(id, name, created_at) VALUES (?, ?, ?)" } }逻辑说明:startup.mode = "initial"是理解这个任务的钥匙。它先对table-names里指定的表做一轮一致性快照,快照完成后再自动切换到 binlog 增量,整个过程对上层应用无感知。job.mode必须设成STREAMING,因为增量阶段是长驻任务,不会自己退出。server-id这里写了一个范围5400-5404,这是给 CDC 读取进程用的伪从库标识,范围要覆盖并行度,避免多个并行实例撞同一个 server-id。
参数说明:database-name是库名,table-names是带库名的表清单,支持正则写法;如果只需要增量不要历史数据,把startup.mode改成latest即可。使用 MySQL-CDC 前需要先确认源库满足三个前置条件,否则任务会卡在“等待 binlog”阶段一动不动:源库my.cnf里打开log-bin、binlog_format=ROW,同步账号有SELECT、REPLICATION SLAVE、REPLICATION CLIENT权限。给权限的 SQL 我放在源码包的 docs 目录里,核心就是两条GRANT语句。
这个模板最容易踩的坑是账号权限遗漏:快照阶段正常,增量阶段永远不输出数据,日志里反复出现权限相关的错误。所以我的习惯是上线前先手动查一遍:
SHOW VARIABLES LIKE 'log_bin'; SHOW VARIABLES LIKE 'binlog_format';log_bin必须是ON,binlog_format必须是ROW。两个条件不满足时,配置写得再对也跑不出增量,这不是 SeaTunnel 的问题,是源库没有开启对应的能力。
3.2 Kafka 到 Console:验证数据链路的最快方式
很多时候你只是想确认“Kafka 里的数据能不能被 SeaTunnel 读出来”,此时最合适的 sink 不是数据库,而是控制台。案例包里的kafka_to_console.conf就是干这个的:
env { parallelism = 1 job.mode = "STREAMING" } source { Kafka { bootstrap.servers = "127.0.0.1:9092" topics = "demo-topic" consumer.group.id = "seatunnel-demo" start.mode = "earliest" format = "json" result_table_name = "kafka_raw" } } sink { Console { source_table_name = "kafka_raw" limit = 100 } }逻辑说明:Kafka source 把 topic 里的消息按format解析成结构化数据,输出到result_table_name;Console sink 把它打到控制台,limit = 100限制最多打印 100 条,防止刷屏。start.mode = "earliest"表示从头消费,适合验证链路;如果只想看新消息,改成latest。
参数说明:bootstrap.servers是 Kafka broker 地址,consumer.group.id决定消费组,topics支持逗号分隔多个 topic。验证完成后,这个任务要用Ctrl+C手动结束,因为STREAMING任务不会自己退出。Console sink 在部分小版本里叫Log,命名以connectors/目录下实际存在的 jar 为准,配置名写错了会直接报插件找不到。
这个模板的价值在于把“数据读取”和“数据写入”两个环节拆开验证。Kafka 消费不到数,先看start.mode和 group id;消费到了但格式不对,检查format是否和消息实际编码一致,这个判断比直接写库再查库快得多。
3.3 批量落地 ClickHouse:fields 显式声明,别让黑匣子猜字段
MySQL 同步到 ClickHouse 是分析型业务最常见的落地场景。ClickHouse 和 MySQL 表结构差异大,类型映射经常出幺蛾子,所以案例包里的mysql_to_clickhouse.conf建议把字段列表显式写清楚:
env { parallelism = 2 job.mode = "BATCH" } source { Jdbc { url = "jdbc:mysql://127.0.0.1:3306/demo_db?serverTimezone=Asia/Shanghai&useSSL=false" user = "root" password = "123456" query = "SELECT id, name, created_at FROM demo_user" result_table_name = "src_user" } } sink { Clickhouse { host = "127.0.0.1" port = 8123 database = "demo_db" table = "demo_user_ch" fields = ["id", "name", "created_at"] username = "default" password = "123456" bulk_size = 20000 } }逻辑说明:Clickhouse sink 的fields数组显式声明写入字段,顺序必须和 source 输出列一致。不要依赖“自动映射”,因为 ClickHouse 与 MySQL 的类型推断规则不完全一致,自动映射翻车时日志报的是类型转换异常,很难一眼定位是哪个字段。bulk_size = 20000控制每次批量写入的行数,这个值对 ClickHouse 的写入性能和合并压力影响很大:太小写入次数多,太大单批次内存占用高。
参数说明:Clickhouse sink 走的是 HTTP 接口,host和port填 clickhouse-server 的地址和 HTTP 端口(默认 8123),username默认是default。不同 SeaTunnel 版本里这个插件名可能写成ClickHouse或Clickhouse,大小写有差异,写之前先看connectors/clickhouse/lib下的 jar 名。
这里有一个经验值:MySQL 到 ClickHouse 的数据迁移,建议先在目标端建好ReplacingMergeTree或MergeTree表,把主键、排序键和分区键定好,再由 SeaTunnel 只负责搬运。否则 SeaTunnel 任务每天跑,ClickHouse 里的表结构一把梭,三个月后查询性能会差到怀疑人生。
3.4 一个任务多 source 多 sink:source_table_name 必须显式声明
同步需求很少是一张表对一张表。案例包里multi_source_multi_sink.conf演示了单任务内多张表并行搬运的写法,核心是每个 sink 都要用source_table_name声明自己消费哪路数据:
env { parallelism = 2 job.mode = "BATCH" } source { Jdbc { url = "jdbc:mysql://127.0.0.1:3306/demo_db" user = "root" password = "123456" query = "SELECT id, name FROM demo_user" result_table_name = "src_user" } Jdbc { url = "jdbc:mysql://127.0.0.1:3306/demo_db" user = "root" password = "123456" query = "SELECT id, amount FROM demo_order" result_table_name = "src_order" } } sink { Jdbc { source_table_name = "src_user" url = "jdbc:mysql://127.0.0.1:3306/demo_db" user = "root" password = "123456" query = "INSERT INTO user_copy(id, name) VALUES (?, ?)" } Jdbc { source_table_name = "src_order" url = "jdbc:mysql://127.0.0.1:3306/demo_db" user = "root" password = "123456" query = "INSERT INTO order_copy(id, amount) VALUES (?, ?)" } }逻辑说明:两个 source 分别产出src_user和src_order两路数据,两个 sink 各自声明消费哪一路。如果没有source_table_name,sink 会默认消费唯一的 source;一旦任务里有多个 source,不声明就会出现“数据被第一个 sink 抢走,第二个 sink 等不到数据”的假死现象。
参数说明:多 source 场景下,每个 source 的result_table_name必须唯一,parallelism是全局配置,会同时作用于所有 source 和 sink。需要把两路数据合并成一路再写入时,不能在 sink 里直接写 join,得先加一个transform段用 join 组件处理,这个玩法源码包里有一个注释版示例,新手可以直接照着改。
4. 配置跑不起来的排查笔记:五个必查的坑
4.1 现象:ClassNotFoundException,或 “No suitable driver found for jdbc:mysql”
配置里 URL、账号都对,任务一启动就报驱动类找不到。原因很简单:SeaTunnel 的 Jdbc connector 只是通过 JDBC 接口连数据库,真正的 JDBC 驱动需要单独放进lib/目录。连接器 jar 和驱动 jar 是两回事,很多人只装了连接器,以为就够了。解决:根据数据库版本下载对应驱动,MySQL 8.x 用mysql-connector-j,拷贝到lib/目录后重启任务,验证命令ls -lh lib | grep mysql,能看到 jar 才算数。这个坑占了配置类报错的三分之一,排查顺序永远是:先看驱动在不在,再看 URL 对不对,最后才怀疑语法。
4.2 现象:CDC 任务快照完成之后,增量阶段一动不动
同步账号权限没问题,binlog 相关的配置也开了,任务快照阶段正常结束,但增量阶段迟迟不吐数据。原因大概率是源库账号缺少REPLICATION SLAVE或REPLICATION CLIENT权限,SeaTunnel 进程在反复请求 binlog 位置但被拒绝。解决:补上两种权限。我一般在测试环境用 root 账号跑通,再收敛成最小权限账号,这样能快速区分“权限问题”和“配置问题”。验证权限的方式是手动执行SHOW MASTER STATUS,如果查询报权限错误,说明 SeaTunnel 也会在同一个位置卡住。
4.3 现象:batch 任务一直跑不结束,日志里没有报错
配置是BATCH模式,启动后任务却不退出。常见原因有两个:一是job.mode被误写成STREAMING,任务常驻等待新数据;二是parallelism设置过高,源库连接池被打满,任务在排队等连接。解决:先确认job.mode是BATCH,再把parallelism降到 2 重跑。判断方法是看日志里有没有JOB_FINISHED,一直没有就按这两条路查,不要在并行度上较劲,先跑通再调优。
4.4 现象:同步成功但目标表数据错位,某列值串到另一列
任务正常结束,行数也对,但数据内容张冠李戴。原因:sink 的query里?占位符顺序和 source 的 SELECT 列顺序不一致,SeaTunnel 的 Jdbc sink 按位置绑定参数,不做列名匹配。解决:把 source 的 SELECT 列顺序和 sink 的 INSERT 列顺序写成完全一致,两边都显式列名。这个坑在 source 表加了新字段后特别容易复发,所以我在源码包的每个 sink 配置里都写了字段注释,改表结构时先改注释再改 SQL。
4.5 现象:中文写入后变成问号或乱码
数据链路通了,但中文全部变成?。原因:连接 URL 里缺characterEncoding=UTF-8,或者目标表字符集不是utf8mb4。解决:URL 统一拼上characterEncoding=UTF-8,目标表建表时用utf8mb4。这里有个细节:utf8mb4才能存 emoji,普通utf8遇到特殊字符一样会丢,这个和 MySQL 本身的语义有关,不是 SeaTunnel 能兜底的。排查时先看源端 SELECT 出来是否正常,再看目标表字符集,最后看 URL 参数,三步走完基本能定位。
5. 让配置更稳的工程化手段:并行度、重试与脚本管理
5.1 并行度不是越大越好:先跑通再慢慢加压
很多人拿到配置第一件事就是调高parallelism,觉得并行度越大速度越快。真实情况是:并行度受限于源库连接数、目标库写入能力、单条数据大小三者的最小值。把并行度从 1 调到 4,速度可能只涨 30%;从 4 调到 16,源库连接池直接被打爆,报Too many connections。SeaTunnel 的env.parallelism是全局并发控制,source 和 sink 都会受影响。
我在源码包里给了一套保守的调参顺序:先用parallelism = 2跑通,记录writeRecordCount和总耗时;然后每次加 2 并观察源库监控,耗时不再下降就停。遇到目标端写入慢的场景,优先把 sink 的批量参数调大,而不是继续加并行度。比如 MySQL sink 的batch_size从 1000 调到 5000,往往比并行度翻倍更有效。并行度配置的作用范围可以用下面这张表概括:
| 参数位置 | 控制范围 | 建议起点 | 调优方向 |
|---|---|---|---|
| env.parallelism | 全局默认并行度 | 2 | 按源库负载逐步上调 |
| source.parallelism | 单个 source 读取并发 | 1 | CDC 场景保持 1,避免 server-id 冲突 |
| sink 的 batch_size / bulk_size | 单批写入行数 | 1000 / 20000 | 目标库写入慢时优先调这个 |
5.2 重试与批量参数:别让瞬时错误杀死任务
同步任务大部分失败不是配置错,而是网络抖动、目标库临时拒绝连接。SeaTunnel 的 Jdbc sink 支持max_retries参数,控制写入失败后的重试次数。实际配置时我会把max_retries设为 3,配合batch_size一起用:批次越小,单次事务时间越短,失败后重试的开销越小。批次设得过大,一个批次里有一条脏数据会导致整个批次回滚,重试几轮都过不去,日志里全是重复的异常堆栈。
设置重试参数时注意区分两个概念:重试次数是“这一批写失败再试几次”,不是“整个任务失败后重启几次”。任务级别的容错要靠调度层处理,比如定时任务用 crontab 拉起,而不是在配置里无限加大max_retries。重试次数设太高,目标库持续不可用时会反复堆积重试请求,反而拖垮恢复中的服务。我的惯例是max_retries = 3,连续失败就停止,让告警暴露问题,而不是让任务在原地空转。
5.3 把配置交给 shell 脚本:日志、备份与回滚
配置稳定之后,下一个目标是让运维可操作。源码包里附带了一个极简的job_manager.sh,作用就三个:按任务名启动、按任务名停止、日志按任务名分文件。脚本核心逻辑只有十几行:
#!/bin/bash # job_manager.sh start|stop job_name SEATUNNEL_HOME=/opt/apache-seatunnel CONF_DIR=$SEATUNNEL_HOME/config/jobs LOG_DIR=$SEATUNNEL_HOME/logs start() { local job=$1 # 每次启动前自动备份当前配置,这是回滚的后悔药 cp "$CONF_DIR/$job.conf" "$CONF_DIR/bak/$(date +%F)_$job.conf" nohup "$SEATUNNEL_HOME/bin/seatunnel.sh" --config "$CONF_DIR/$job.conf" \ >> "$LOG_DIR/$job.log" 2>&1 & echo "started: $job, pid: $!" } stop() { pkill -f "$CONF_DIR/$1.conf" } case "$1" in start) start "$2" ;; stop) stop "$2" ;; esac逻辑说明:启动前先把当前配置备份到bak/目录,文件名为日期加任务名,这样任何一次改配置导致的任务失败,都能回滚到前一天能跑的版本。日志按任务名落到独立文件,排查问题时直接grep -E "JOB_FINISHED|writeRecordCount" logs/xxx.log,不用在一大坨混合日志里翻找。
这套脚本的定位是本地和测试环境。生产环境如果任务多、要求高可用,建议切到 Zeta 集群模式,通过作业 API 提交和管理任务,脚本方案应付不了故障转移。但从单机到集群的迁移过程中,这个脚本的“备份再启动”思路可以原样带过去:任何配置变更都要留一份可回滚的副本,这是我在生产环境交过学费之后养成的习惯。
6. 一个具体技巧:用 variable 把配置变成多环境模板
配置里写死连接信息是同步任务最常见的隐患。开发、测试、生产三套库,同一个配置文件要维护三份,每次发版都要人肉改。SeaTunnel 支持用--variable参数在运行时替换配置里的${变量},把同一个配置文件变成多环境通用模板。
配置文件里这样写占位符:
source { Jdbc { url = "jdbc:mysql://${DB_HOST}:${DB_PORT}/${DB_NAME}?serverTimezone=Asia/Shanghai&useSSL=false" user = "${DB_USER}" password = "${DB_PASSWORD}" query = "SELECT id, name FROM ${SOURCE_TABLE}" } }启动时通过命令行传入变量值:
bin/seatunnel.sh --config config/user_sync_template.conf \ --variable "DB_HOST=127.0.0.1" \ --variable "DB_PORT=3306" \ --variable "DB_NAME=demo_db" \ --variable "DB_USER=root" \ --variable "DB_PASSWORD=123456" \ --variable "SOURCE_TABLE=demo_user"变量名我习惯用大写加下划线,避免和配置文件里的特殊字符解析冲突。不同小版本对--variable的支持略有差异,用之前先跑一遍bin/seatunnel.sh --help确认参数名;如果版本不支持变量替换,还有一个土办法:启动脚本里用sed根据环境文件生成临时 conf,任务跑完即删。两种方式我都用过,--variable更规范,sed方案兼容老版本。
我有一年在生产环境同步用户表,连接信息写死在三个配置文件里。某天目标库切换,我漏改了 sink 段里的库名,凌晨两点同步任务照常“成功”,数据全进了旧库,第二天对账才发现。从那以后,我每次部署同步任务都强制走一遍变量替换,把环境差异全部收口到环境文件里,配置模板半个字都不动,改环境只改变量。希望帮到你。
本文还有配套的精品资源,点击获取