SeaTunnel 多表 CDC 实战:一个作业同步 MySQL 多张表并按表自动路由到 PostgreSQL
2026/9/18 4:24:24 网站建设 项目流程

SeaTunnel 多表 CDC 实战:一个作业同步 MySQL 多张表并按表自动路由到 PostgreSQL

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

本文是一份基于 Apache SeaTunnel 的「多表 CDC」实战配方(Recipe):通过单个流式作业,用 MySQL-CDC Source 的正则table-pattern捕获多张上游表的变化,再用 JDBC Sink 的${table_name}占位符将每个源表自动路由到 PostgreSQL 中各自独立的目标表(如st_ordersst_customersst_products),并借助generate_sink_sql = true自动建表、自动生成 UPSERT 语句。读完本文你将掌握:CDC 前置环境(binlog、用户权限、驱动)的搭建方法、最小可运行的多表 CDC 配置、端到端验证步骤,以及多表路由在 SeaTunnel 内部的源码级实现原理。

适用场景

当你有以下需求时,本配方是标准解法:

  • 单作业多表:不希望为每张表各起一个同步作业,而是用一个作业批量捕获数十上百张表;
  • 一致的快照起点:多张表需要从同一时间点开始全量快照,之后无缝切换到 binlog 增量;
  • 按表自动路由:每张源表的数据自动落到对应的目标表,而不是被混写进同一张表;
  • 流式持续同步:作业常驻运行,MySQL 侧的新增、更新持续流入目标库。

多表同步的架构目标与设计取舍可进一步阅读 多表同步架构文档,本文先聚焦如何把它跑起来。

前置条件

1. 已完成第一个作业

先确保你能跑通 SeaTunnel 的基本作业流程,参见 运行你的第一个作业。

2. 安装本配方所需的插件

按照 部署文档 > 下载连接器插件 的方式安装插件,并在config/plugin_config中只保留以下两个插件:

--seatunnel-connectors-- connector-cdc-mysql connector-jdbc --end--

执行安装并确认插件就位:

cd "${SEATUNNEL_HOME}" sh bin/install-plugin.sh ls connectors | rg 'connector-(cdc-mysql|jdbc)'

3. 放置 JDBC 驱动(SeaTunnel Zeta 引擎)

如果你使用 SeaTunnel Zeta 引擎,需要同时把 MySQL 与 PostgreSQL 的 JDBC 驱动放入${SEATUNNEL_HOME}/lib,然后确认可见:

ls "${SEATUNNEL_HOME}/lib" | rg 'mysql-connector|postgresql'

SeaTunnel 不捆绑所有 JDBC 驱动,驱动 JAR 需自行下载;Zeta 引擎统一从${SEATUNNEL_HOME}/lib加载驱动,放置后需重启 SeaTunnel 进程。如果使用 Spark/Flink 引擎,则需放到${SEATUNNEL_HOME}/plugins/Jdbc/lib/(详见 JDBC Sink 文档)。

4. 准备 MySQL 源表

每张上游表都应具备稳定的主键,因为本配方会把 CDC 变更自动路由到下游的 upsert 目标表(主键同时是路由与去重的依据):

CREATE DATABASE IF NOT EXISTS inventory; CREATE TABLE IF NOT EXISTS inventory.orders ( id BIGINT PRIMARY KEY, order_status VARCHAR(32), updated_at TIMESTAMP DEFAULT CURRENT_TIMESTAMP ON UPDATE CURRENT_TIMESTAMP ); CREATE TABLE IF NOT EXISTS inventory.customers ( id BIGINT PRIMARY KEY, customer_name VARCHAR(64), city VARCHAR(64), updated_at TIMESTAMP DEFAULT CURRENT_TIMESTAMP ON UPDATE CURRENT_TIMESTAMP ); CREATE TABLE IF NOT EXISTS inventory.products ( id BIGINT PRIMARY KEY, product_name VARCHAR(64), unit_price DECIMAL(10, 2), updated_at TIMESTAMP DEFAULT CURRENT_TIMESTAMP ON UPDATE CURRENT_TIMESTAMP ); INSERT INTO inventory.orders (id, order_status, updated_at) VALUES (2001, 'CREATED', NOW()); INSERT INTO inventory.customers (id, customer_name, city, updated_at) VALUES (3001, 'Alice', 'Shanghai', NOW()); INSERT INTO inventory.products (id, product_name, unit_price, updated_at) VALUES (4001, 'Keyboard', 99.00, NOW());

关于无主键表:MySQL CDC 默认期望源表有主键。若某张表没有主键但存在唯一列,可通过table-names-config.primaryKeys指定自定义主键;如果完全没有稳定的唯一键,UPDATE/DELETE 事件将无法安全地在下游应用(详见 MySQL CDC 文档)。

5. 创建 MySQL CDC 用户并授权

CDC 用户需要SELECT(读快照)、RELOADSHOW DATABASES,以及REPLICATION SLAVE, REPLICATION CLIENT(读 binlog):

CREATE USER IF NOT EXISTS 'st_user_source'@'%' IDENTIFIED BY 'mysqlpw'; GRANT SELECT, RELOAD, SHOW DATABASES, REPLICATION SLAVE, REPLICATION CLIENT ON *.* TO 'st_user_source'@'%'; FLUSH PRIVILEGES;

6. 确认 MySQL binlog 配置

SHOW VARIABLES WHERE variable_name IN ('log_bin', 'binlog_format', 'binlog_row_image');

期望值为log_bin = ONbinlog_format = ROWbinlog_row_image = FULL。若未开启,需在 MySQL 配置文件中启用:

[mysqld] server-id = 223344 log_bin = mysql-bin expire_logs_days = 10 binlog_format = row binlog_row_image = FULL

并重启 MySQL。注意:当初始快照较大时,还需适当调大interactive_timeoutwait_timeout,防止快照期间连接超时(MySQL CDC 文档 中给出了更完整的 binlog 与会话超时说明)。

7. 准备 PostgreSQL 目标库与写入用户

CREATE USER st_user_sink WITH PASSWORD 'pgpw'; CREATE DATABASE sync_demo; GRANT ALL PRIVILEGES ON DATABASE sync_demo TO st_user_sink;

重新连接sync_demo后,授予publicschema 上的建表权限:

GRANT USAGE, CREATE ON SCHEMA public TO st_user_sink;

本配方使用了generate_sink_sql = true,因此 SeaTunnel 会在首次运行时自动创建public.st_orderspublic.st_customers等目标表,前提是 sink 用户对publicschema 拥有CREATE权限。

最小配置

下面这份配置通过一个table-pattern读取 MySQL 的多张表,并写入 PostgreSQL 中名为st_<上游表名>的目标表:

env { parallelism = 1 job.mode = "STREAMING" checkpoint.interval = 5000 } source { MySQL-CDC { plugin_output = "mysql_multi" startup.mode = "initial" server-id = 5652 username = "st_user_source" password = "mysqlpw" database-pattern = "inventory" table-pattern = "inventory\\.(orders|customers|products)" url = "jdbc:mysql://mysql:3306/inventory" } } sink { Jdbc { plugin_input = "mysql_multi" driver = "org.postgresql.Driver" url = "jdbc:postgresql://postgresql:5432/sync_demo" username = "st_user_sink" password = "pgpw" generate_sink_sql = true database = "sync_demo" table = "public.st_${table_name}" primary_keys = ["${primary_key}"] schema_save_mode = "CREATE_SCHEMA_WHEN_NOT_EXIST" data_save_mode = "APPEND_DATA" } }

关键参数逐项解读

env 段

参数说明
parallelism1本配方用单并行度演示;多表场景可调大并行度以提升吞吐
job.mode"STREAMING"流式作业,快照完成后持续消费 binlog
checkpoint.interval5000每 5 秒做一次 checkpoint,决定下游可见性与故障恢复粒度

source 段(MySQL-CDC)

参数说明
plugin_output"mysql_multi"为数据流命名,下游plugin_input须与之对应
startup.mode"initial"启动时先做全量快照,再无缝切换增量;其他可选值:earliestlatestspecifictimestamp
server-id5652该 CDC 读取器在 MySQL 集群中的唯一 ID,也可配区间如5652-5660;多读取器/多表并行时区间要足够大,重复会导致 MySQL 踢掉其中一个客户端
database-pattern"inventory"要捕获的库名正则,例如database_prefix.*
table-pattern"inventory\\.(orders\|customers\|products)"要捕获的表名正则,匹配名包含库名;table-patterntable-names互斥,二选一
urljdbc:mysql://mysql:3306/inventoryJDBC 连接地址

其中table-pattern是「一个作业捕获多表」的关键:SeaTunnel 会把正则命中的每张表各自解析成独立的数据分片与元数据(表结构、主键),而不是把多张表混成一个 schema。对应的选项定义与互斥校验可查看源码 MySqlIncrementalSourceFactory.java 中对TABLE_NAMES/TABLE_PATTERN的解析逻辑。

sink 段(Jdbc)

参数说明
plugin_input"mysql_multi"对应上游plugin_output
driver/urlPostgreSQL 驱动与地址本配方向 PostgreSQL 写入
generate_sink_sqltrue由 SeaTunnel 依据上游 schema 与行类型(INSERT/UPDATE/DELETE)自动生成 SQL;若为false则必须手写query
database"sync_demo"目标库名
table"public.st_${table_name}"路由核心${table_name}占位符会被替换为上游记录携带的真实表名,于是orderspublic.st_orderscustomerspublic.st_customers
primary_keys["${primary_key}"]同样支持占位符,从上游元数据继承每张表的主键,用于生成数据库原生的 UPSERT / UPDATE / DELETE
schema_save_mode"CREATE_SCHEMA_WHEN_NOT_EXIST"目标表不存在时自动创建(RECREATE_SCHEMA会先删后建、ERROR_WHEN_SCHEMA_NOT_EXIST缺失即报错、IGNORE跳过建表逻辑)
data_save_mode"APPEND_DATA"保留已有数据继续追加(DROP_DATA会清空、ERROR_WHEN_DATA_EXISTS有数据即报错)

关于${table_name}占位符:JDBC Sink 的table参数支持${table_name}${schema_name}变量,${schema_name}会被替换为目标侧 schema 名,${table_name}会被替换为目标侧表名。示例写法如test_${schema_name}_${table_name}_testpublic.${table_name}_test等(详见 JDBC Sink 文档)。对带 schema 概念的数据库(如 PostgreSQL、Oracle、SQL Server),table需写成xxx.xxx形式。

运行作业

将配置保存为config/multi-table-cdc.conf,然后以本地模式启动 SeaTunnel:

cd "${SEATUNNEL_HOME}" ./bin/seatunnel.sh --config ./config/multi-table-cdc.conf -m local

因为这是一个流式 CDC 管道,请保持作业运行,同时到 MySQL 侧制造新的变更来验证增量同步。

验证结果

第一步:确认建表

等初始快照阶段结束后,在 PostgreSQL 中查询 SeaTunnel 自动创建的目标表:

SELECT table_name FROM information_schema.tables WHERE table_schema = 'public' AND table_name LIKE 'st_%' ORDER BY table_name;

预期能看到st_ordersst_customersst_products三张表。

第二步:在 MySQL 各源表上制造变更

INSERT INTO inventory.orders (id, order_status, updated_at) VALUES (2002, 'PAID', NOW()); UPDATE inventory.customers SET city = 'Hangzhou', updated_at = NOW() WHERE id = 3001; INSERT INTO inventory.products (id, product_name, unit_price, updated_at) VALUES (4002, 'Mouse', 59.00, NOW());

第三步:确认每张下游表只收到自己的数据

SELECT id, order_status FROM public.st_orders ORDER BY id; SELECT id, customer_name, city FROM public.st_customers ORDER BY id; SELECT id, product_name, unit_price FROM public.st_products ORDER BY id;

st_orders应包含 2001/2002 两条订单,st_customers中 Alice 的城市应已变为 Hangzhou(UPSERT 生效),st_products应包含 4001/4002 两条产品。如果每个上游表都被路由到各自的目标表且变更持续流动,说明多表 CDC 管道工作正常。

底层原理:多表路由是怎么实现的

理解了「怎么配」,再看「为什么这样就能路由」。整条链路由三层机制协作完成:

1. Source 侧:正则展开为多张表的独立元数据

table-pattern命中多张表后,MySQL CDC Source 会为每张表生成独立的 catalog 元数据(CatalogTable)与读取分片(split)。每张表拥有自己的 schema 与主键信息,为下游按表建表和按表生成 SQL 提供了基础。全量快照阶段完成后,Source 会自动切换到从快照开始时记录的 binlog 位点继续消费增量,保证切换期间不丢事件(MySQL CDC 文档 的 FAQ 有明确说明)。

server-id在源码层面被解析为ServerIdRange,每个并行子任务会从区间中分得一个唯一 ID 并写入 Debezium 配置database.server.id,避免多读取器在同一 MySQL 集群上发生 server-id 冲突(见 MySqlSourceConfigFactory.java)。

2. 传输侧:每行数据都携带 tableId

多表管道中,每一行SeaTunnelRow都会携带一个tableId(即TablePath的序列化形式:databaseName.schemaName.tableName)。TablePath是记录归属的唯一标识,定义于 TablePath.java。正是这个 tableId 让下游无需猜测「这行数据属于哪张表」。

3. Sink 侧:MultiTableSink 按表分桶写入

JDBC Sink 在接收到多表输入时,会被包装为框架层的MultiTableSink。它会为每张表创建独立的子 Sink(writer),并在运行时根据行的tableId把记录路由到对应子 writer。核心实现见 MultiTableSink.java 与 MultiTableSinkWriter.java。

路由策略分两种情况:

  • 有主键时(哈希路由):对主键字段值取哈希并映射到队列下标,保证同一主键的记录永远进入同一个队列,从而保持该键内的写入顺序:

    int index = (object.hashCode() & Integer.MAX_VALUE) % blockingQueues.size();

    源码中刻意用& Integer.MAX_VALUE清掉符号位而非Math.abs,因为Math.abs(Integer.MIN_VALUE)仍返回负数,会导致下标越界。

  • 无主键时(随机路由):在各队列间随机分布以均衡负载,但不保证同一键的顺序。

每个队列对应一个 replica(副本 writer),可通过multi_table_sink_replica配置每个表的并发 writer 数量以提升吞吐。schema 变更(如新增列)则以 barrier 形式广播到所有队列,确保每张表的行流在同一个位置应用变更、不乱序。

4.${table_name}占位符如何生效

子 Sink 写入时,table参数中的${table_name}${schema_name}会被替换为该行所属TablePath中的真实表名、schema 名;primary_keys中的${primary_key}则替换为该表元数据中的主键列。配合generate_sink_sql = true,SeaTunnel 即可为每张表生成专属的CREATE TABLEINSERT ... ON CONFLICT ... DO UPDATE(PostgreSQL 原生 UPSERT)语句。这就是「每张源表自动路由到各自目标表」的完整闭环。

常见陷阱

以下是本配方最容易踩的坑,逐一对照排查:

  • table-pattern中的正则转义错误:在 HOCON 里,.表示任意字符,匹配字面点号时必须写成\\.。例如inventory\\.(orders|customers|products),写错会导致表匹配不到或误匹配。
  • MySQL binlog 或 CDC 用户权限不完整:表现是作业能读完快照,但无法继续读增量。核对binlog_format = ROWbinlog_row_image = FULL,并确认用户具备REPLICATION SLAVE, REPLICATION CLIENT
  • 未配置占位符路由:如果 sink 的table写死成单张表名(没有${table_name}),多张源表会被意外混写进同一张目标表。
  • PostgreSQL sink 用户缺少CREATE权限:能连接数据库,但无法在publicschema 上建表。需执行GRANT USAGE, CREATE ON SCHEMA public TO st_user_sink;
  • 上游表没有主键却按 upsert 语义配置:没有稳定主键时,UPSERT/UPDATE/DELETE 无法安全应用。要么给表加主键,要么用table-names-config.primaryKeys指定唯一列。
  • schema/table 占位符用错了位置:不同数据库对 schema 的语义不同(如 MySQL 无独立 schema、PostgreSQL 有public),${schema_name}${table_name}的用法需与目标库的命名规则匹配。

相关文档

  • MySQL CDC 连接器文档:table-pattern/table-names互斥规则、server-id区间、startup.mode全量增量切换、binlog 与用户权限的完整说明
  • JDBC Sink 连接器文档:generate_sink_sql两种写入模式、${table_name}/${schema_name}占位符规则、schema_save_mode/data_save_mode取值、UPSERT 行为
  • 多表同步架构文档:TablePathMultiTableSink、replica 副本机制、schema 变更路由与故障隔离策略的源码级剖析

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

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

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

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

立即咨询