SeaTunnel 数据集成上手:从本地跑通到集群稳跑,一篇讲清楚
2026/9/18 17:05:01 网站建设 项目流程

SeaTunnel 数据集成上手:从本地跑通到集群稳跑,一篇讲清楚

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

公司里最常见的数据集成需求,往往不是简单的"搬一张表":从 MySQL 拉订单,只留已支付的,字段清洗一下,再写进分析库;或者把 Kafka 里的订单事件持续落到湖表供查询。这类活儿以前要么靠自研脚本,要么上重量级流平台。今天聊的 SeaTunnel 数据集成工具,就是一个专门干这件事的开源方案:你不用写代码,写一份配置文件描述"数据从哪来、变成什么样、到哪去",然后一条命令提交任务。下面按"先跑起来、再看懂、最后上生产"的顺序走一遍。

三步在本地跑通第一个 SeaTunnel 任务

先不谈架构,把环境搭起来。要求很朴素:一台装了 JDK 8 或 11 的机器,4GB 以上内存。

# 1. 下载官方二进制包,解压 tar -xzf apache-seatunnel-2.3.13-bin.tar.gz cd apache-seatunnel-2.3.13 # 2. 按 config/plugin_config 的清单安装连接器插件 sh bin/install-plugin.sh

plugin_config是一个插件清单文件,默认只装了最少的几个。为什么单独装插件而不是全部?因为 SeaTunnel 支持一百多种连接器,全装会占掉不少磁盘和类加载时间,按需装更干净。

跑通验证用一个自带示例任务就够了,它不需要任何外部数据库:

./bin/seatunnel.sh --config ./config/v2.batch.config.template -m local

-m local表示任务在本机进程内执行,不依赖任何远程集群,适合验证安装。任务跑完,控制台会打印 16 行 FakeSource 造的假数据,末尾还有一段 Job Statistic 汇总:读多少条、写多少条、失败多少条。看到Total Failed Count = 0,这条链路就算通了。

看懂配置:数据从哪来、经过什么、到哪去

跑通之后,打开那份模板配置,整个文件其实就是四块:

env { job.mode = "BATCH" parallelism = 1 } source { FakeSource { plugin_output = "fake" } } transform { FieldMapper { plugin_input = "fake" plugin_output = "fake1" } } sink { Console { plugin_input = "fake1" } }
  • env管"怎么跑":job.modeBATCH(跑完即退出)或STREAMING(常驻消费),parallelism是作业并行度。
  • source管数据从哪来,sink管写到哪去,中间的transform是可选的,链路简单时整段可以删掉。
  • 各块之间靠plugin_output/plugin_input这两个名字对上,相当于给每段数据流起了个变量名,下游指名引用。只有一个上游时可以省,但链路一多,显式命名能省掉大量排查时间。

记住这三件事——env / source / sink 三段式、流批一个模式参数切换、插件之间用名字串联——以后看任何 SeaTunnel 数据集成配置都不会迷路。

用例一:批量迁移 MySQL 订单,顺手做清洗

真实业务里最典型的场景:把 MySQL 里的订单迁到 PostgreSQL,只要已支付的,顺便统一字段格式。整条链路的配置骨架是这样(完整可运行版本在仓库的 jdbc-to-jdbc 教程 里,含建表 SQL 和验证步骤):

env { parallelism = 2 job.mode = "BATCH" } source { Jdbc { plugin_output = "mysql_orders" driver = "com.mysql.cj.jdbc.Driver" url = "jdbc:mysql://mysql-host:3306/source_db" query = "SELECT id, customer_name, amount, status FROM orders" } } transform { Sql { plugin_input = "mysql_orders" plugin_output = "paid_orders" query = """ SELECT id, UPPER(customer_name), CAST(amount AS DECIMAL(12,2)), 'MYSQL' AS source_system FROM dual WHERE status = 'PAID' """ } } sink { Jdbc { plugin_input = "paid_orders" driver = "org.postgresql.Driver" url = "jdbc:postgresql://pg-host:5432/target_db" query = "INSERT INTO public.paid_orders (id, customer_name, amount, source_system) VALUES (?, ?, ?, ?)" } }

几个参数为什么这么写:

  • parallelism = 2:单机批任务给 2 个并行度,既摊开了读库吞吐,又不至于把 MySQL 的连接和 CPU 打满。数据量大或跨机房时可再加,加之前先看源库负载。
  • 中间那个Sql转换:过滤、改名、补来源系统字段都在一条 SQL 里做完,比串好几个 transform 插件直观,字段顺序改动时只需要看一处。
  • 跑之前把 MySQL 和 PostgreSQL 的驱动 jar 放进lib/目录,缺驱动是最常见的"启动即报错"原因。
./bin/seatunnel.sh --config ./job/mysql-to-pg.conf -m local

跑完去 PostgreSQL 查目标表,行数对得上、未支付订单确实没进来,这个用例就闭环了。

用例二:把 Kafka 订单流持续落到 Iceberg 表

流式场景的思路完全不同:任务常驻,靠 checkpoint 保证不丢不重。一条从 Kafka 到 Iceberg 的最小配置(完整版见 kafka-to-iceberg 教程):

env { job.mode = "STREAMING" checkpoint.interval = 5000 } source { Kafka { plugin_output = "orders_kafka" topic = "orders" bootstrap.servers = "kafka:9092" consumer.group = "seatunnel-orders" format = "json" } } sink { Iceberg { plugin_input = "orders_kafka" catalog_name = "seatunnel_demo" table = "orders" iceberg.table.primary-keys = "id" iceberg.table.upsert-mode-enabled = true } }
  • job.mode = "STREAMING"且开了checkpoint.interval = 5000:流作业不开 checkpoint,任务重启后的消费位点和落表一致性就没有保障,所以这一项建议直接写进模板。5 秒是"状态保存频率"和"对任务吞吐的额外开销"之间的常见折中,延迟要求更高可以调到 1~2 秒,但 checkpoint 越频繁,元数据写入压力越大。
  • iceberg.table.primary-keys = "id"upsert-mode-enabled:订单事件同一条id可能多次更新,按主键 upsert 而不是无脑追加,表里才是最新状态。
./bin/seatunnel.sh --config ./job/kafka-to-iceberg.conf -m local

Kafka source 按分区切分并行任务,多分区 topic 上并行度可以拉到和分区数一个量级;单分区 topic 加并行度没用,数据本身是串行的。

上生产:集群部署、监控与排障

单机扛不住或者想要任务故障自动恢复时,再谈集群。

部署。核心是config/hazelcast.yaml

hazelcast: cluster-name: seatunnel network: join: tcp-ip: enabled: true member-list: - 192.168.1.100 - 192.168.1.101 - 192.168.1.102

member-list里必须列全集群所有节点,包括 master 自己——Hazelcast 靠这份静态清单互相发现,漏一个节点就会出现"看起来起来了实际在两个脑裂的小集群里"的问题。改完后每个节点执行:

sh bin/seatunnel-cluster.sh -r master # 控制节点 sh bin/seatunnel-cluster.sh -r worker # 执行节点

提交任务时把-m local换成-m cluster即可,作业配置本身一行不用改。

监控config/seatunnel.yamlhttp.port默认是 8080,浏览器打开就能看到作业列表、每个 source/sink 的读写吞吐、运行时长;作业异常时,logs/seatunnel.log是第一现场,JVM 堆转储路径在config/jvm_options里配的是/tmp/seatunnel/dump/,OOM 崩溃后不用现场抓,直接看 dump。

排障三板斧。一是连接器没装:ls connectors/里没有对应目录,回去跑install-plugin.sh,它按plugin_config清单工作;二是驱动缺失:JDBC 类任务报ClassNotFoundException,九成是lib/里没有数据库驱动 jar;三是内存:批量任务 OOM 时优先调jvm_options里的-Xmx并配合env { parallelism }把单任务数据量摊薄,比一味加堆更稳,堆过大反而拖长 GC 停顿。

多团队共用集群时,可以用 tag 把资源圈成几个池子,各团队的任务互不抢资源:

接下来做这三件事

  1. 把用例一里的假主机换成你环境里真实的源库,跑一次端到端小表同步,确认权限、驱动、网络三件事都通。
  2. checkpoint.intervalparallelism建一个团队内的配置模板,新任务从模板改而不是从零写,减少"流作业忘了开 checkpoint"这类低级事故。
  3. 想搞清配置细节和更多链路(MySQL CDC、多表同步等),看仓库的 作业配置指南 和 recipes 目录,每条链路都有可直接抄的最小配置。

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

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

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

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

立即咨询