- 数据集成
- ETL
- 大数据
- 批处理
- 流处理
- 变更数据捕获
【免费下载链接】seatunnel
SeaTunnel is a multimodal, high-performance, distributed, massive data integration tool.
本文围绕 SeaTunnel 官方场景教程中一条已经在 Docker 环境里完成端到端验证的链路展开:使用MySQL-CDC读取shop.orders表的历史快照与后续 binlog 增量,经Metadatatransform 暴露库名、表名、变更类型,再经Sqltransform 完成字段重命名与业务语义整形,最后由Kafkasink 以 JSON 格式输出 payload,并把部分元数据写入 Kafka 消息头。读完本文,你将掌握这条链路的完整配置、运行前置条件、逐条断言结果,以及kafka_headers_fields同时影响 header 与 payload 的底层行为。
链路概览:一条被实测断言过的四段流水线
这篇教程只讲一条已经在 Docker 里完成端到端验证的链路。验证时间是 2026 年 7 月 16 日,文中保留的配置、输入数据和结果,都是这次实测里真正跑过并断言过的内容。
这条已验证链路是:
MySQL-CDC读取shop.orders的快照和后续 binlogMetadata暴露库名、表名和变更类型Sql重命名字段并补充业务字段Kafka输出 JSON payload,同时把部分元数据写入消息头
与它对应的自动化验证位于仓库的 Kafka 连接器 E2E 测试中:KafkaRecipeIT.java(其类注释明确写着 "Validates the documented MySQL CDC to Kafka recipe with metadata enrichment and SQL field shaping")。该测试通过@DisabledOnContainer限定只在 Zeta 引擎上运行,以匹配 getting-started 文档路径,说明这条 recipe 面向 SeaTunnel Zeta 引擎验证。测试用的作业配置与文档一致,见 mysqlcdc_to_kafka_with_transforms.conf。
这次验证环境具备的前置条件
这条 Docker 端到端验证在任务启动前,已经具备了下面这些条件:
- MySQL 使用了
docker/server-gtids/my.cnf里的 GTID / binlog 配置。测试代码中通过MySqlContainer.withConfigurationOverride("docker/server-gtids/my.cnf")将这份覆盖配置挂入 MySQL 容器,启用 binlog、binlog_format = ROW、binlog_row_image = FULL与 GTID 模式(该覆盖文件位于 MySQL-CDC 连接器的测试资源中)。 MySQL-CDC插件的lib目录里已经放入了 MySQL JDBC driver JAR。E2E 测试通过DependencyJar.of(Driver.class).copyTo(container, "/tmp/seatunnel/plugins/MySQL-CDC/lib")把com.mysql.cj.jdbc.Driver注入到插件的lib目录,这也是 Docker 场景下为 CDC 插件补充驱动的标准做法。- 预先创建了名为
st_user_source的 CDC 用户,并授予了SELECT、RELOAD、SHOW DATABASES、REPLICATION SLAVE、REPLICATION CLIENT、LOCK TABLES权限。测试中由 root 用户执行CREATE USER与GRANT ... ON *.*完成,权限清单与 MySQL-CDC 文档 中的要求一致。 - SeaTunnel 任务启动前,Kafka topic
recipe_mysql_orders已经提前创建好。测试用 KafkaAdminClient以 1 个分区、副本因子 1 创建该 topic;Kafka 容器使用confluentinc/cp-kafka:7.0.9镜像,并在 Docker 网络中以别名kafkaCluster提供访问。
如果你的生产环境要复现这条链路,除上述四点外,还需要确保 MySQL 源库开启了 binlog(log_bin=ON、binlog_format=ROW、binlog_row_image=FULL),并避免多个 CDC 任务使用重叠的server-id范围——重复的server-id会导致 MySQL 主动断开其中一个客户端连接。
已验证的源表数据
E2E 测试先创建了下面这张表,并插入两条初始数据:
CREATE TABLE orders ( id BIGINT NOT NULL PRIMARY KEY, order_no VARCHAR(64) NOT NULL, user_id BIGINT NOT NULL, status INT NOT NULL, amount DECIMAL(10, 2) NOT NULL ); INSERT INTO orders (id, order_no, user_id, status, amount) VALUES (1001, 'ORD-1001', 501, 0, 19.99), (1002, 'ORD-1002', 502, 1, 29.99);任务启动后,测试又执行了这两条增量变更:
UPDATE shop.orders SET status = 2, amount = 39.99 WHERE id = 1001; INSERT INTO shop.orders (id, order_no, user_id, status, amount) VALUES (1003, 'ORD-1003', 503, 0, 59.99);也就是说,测试分别覆盖了快照阶段(两条初始行)、**binlog 增量阶段(一条 UPDATE + 一条 INSERT)**两类数据形态,从而能同时验证"全量 + 增量"的 CDC 完整生命周期。
Docker 实测通过的完整配置
下面这份配置就是 Docker 测试环境里真实跑通过的那份作业配置。里面的主机名(mysql_cdc_e2e、kafkaCluster)是当时测试网络里的服务别名,替换为你自己的地址即可。
env { parallelism = 1 job.mode = "STREAMING" } source { MySQL-CDC { plugin_output = "mysql_orders_raw" url = "jdbc:mysql://mysql_cdc_e2e:3306/shop" username = "st_user_source" password = "mysqlpw" server-id = 5601-5604 table-names = ["shop.orders"] startup.mode = "initial" schema-changes.enabled = false } } transform { Metadata { plugin_input = "mysql_orders_raw" plugin_output = "mysql_orders_with_meta" metadata_fields { Database = source_database Table = source_table RowKind = change_type } } Sql { plugin_input = "mysql_orders_with_meta" plugin_output = "kafka_orders" query = "select id as order_id, order_no, user_id, amount, case when status = 0 then 'CREATED' when status = 1 then 'PAID' when status = 2 then 'SHIPPED' else 'OTHER' end as status_name, source_database, source_table, change_type, CONCAT(source_database, '.', source_table) as source_name, 'mysql_cdc' as sync_source from dual where id is not null" } } sink { Kafka { plugin_input = "kafka_orders" bootstrap.servers = "kafkaCluster:9092" topic = "recipe_mysql_orders" format = json partition_key_fields = ["order_id"] kafka_headers_fields = ["source_database", "source_table", "change_type"] } }env 段:流式作业与并行度
job.mode = "STREAMING"使作业常驻运行,持续消费 binlog 增量;parallelism = 1在单表小数据量验证场景下足够。注意这里未显式配置checkpoint.interval,CDC 增量读取的进度推进依赖 checkpoint 机制,生产环境建议按需补充。
source 段:MySQL-CDC 关键参数
table-names = ["shop.orders"]:表名必须带库名前缀。table-names与table-pattern二选一配置(详见 MySQL-CDC 文档)。server-id = 5601-5604:CDC 读取器使用的数字 ID 范围。每个 ID 在 MySQL 集群中必须唯一;当任务有多个读取并发或并行读取多张表时,需要配置足够大的 ID 范围。未配置时 SeaTunnel 会随机生成 ID,但生产环境建议显式配置。文档 FAQ 也提示:建议为每个任务分配独立范围(如一个任务5400-5600、另一个5601-5800)。startup.mode = "initial":启动时先同步历史数据(一致性快照),再自动切换到 binlog 增量。有效枚举值为initial、earliest、latest、specific、timestamp。schema-changes.enabled = false:模式演进默认关闭,此时 DDL 变更不会向下方传递。若需传播add column、drop column、rename column、modify column等结构变更,需显式开启(本文链路刻意关闭,聚焦行级数据整形)。
transform 段:元数据提取与 SQL 整形
Metadata只做"把隐藏的行级元数据变成普通字段",不改变原有数据字段。配置语法metadata_fields { Database = source_database }的含义是:将元数据 KeyDatabase投影为输出字段source_database。根据 Metadata 文档,本链路用到的三个 Key 含义如下:
| 元数据 Key | 输出类型 | 说明 |
|---|---|---|
Database | string | 数据所属的数据库名称(所有连接器均提供) |
Table | string | 数据所属的表名称(所有连接器均提供) |
RowKind | string | 行的变更类型,值为+I(插入)、-U(更新前)、+U(更新后)、-D(删除) |
Sqltransform 在元数据字段之上统一做业务整形:id as order_id完成字段重命名,CASE WHEN status = ...把数字状态码翻译成业务语义字符串,CONCAT(source_database, '.', source_table)拼出source_name,'mysql_cdc' as sync_source注入固定的同步来源标识。from dual where id is not null是 SQL 转换引擎的常规写法,用于以无源表的方式执行表达式投影。
sink 段:Kafka 输出与消息头
format = json:payload 默认使用 JSON 格式。Kafka sink 还支持text、canal_json、debezium_json、compatible_debezium_json、ogg_json、maxwell_json、avro、protobuf、native等格式(详见 Kafka sink 文档)。partition_key_fields = ["order_id"]:配置字段用作 Kafka 消息的 key,保证相同order_id的记录落入同一分区,从而在消费侧保持单 key 顺序。若未配置,SeaTunnel 以null作为消息 key 发送,Kafka 会按轮询策略分散到各分区,适合负载均衡但不适合按业务 key 保序的场景。kafka_headers_fields = ["source_database", "source_table", "change_type"]:把这三个字段写进 Kafka 消息头,字段值会被转换为字符串作为 header 值。
这次端到端验证实际证明了什么
Docker E2E 测试对下面这些结果做了断言(对应 KafkaRecipeIT.java 中的Assertions逻辑):
- 快照阶段成功把初始 MySQL 数据写进了 Kafka(断言 topic 中至少包含 2 条记录)。
- 订单
1001的 payload 中包含order_id、status_name、source_name、sync_source,且status_name = CREATED、source_name = shop.orders、sync_source = mysql_cdc。 - 在快照阶段
1001这条记录上,Kafka 消息头里包含source_database=shop、source_table=orders、change_type=+I。 - 在同一条快照记录上,配置了
kafka_headers_fields之后,这三个字段不会再留在 JSON payload 里(payload.has("source_database")等断言为false)。 - 执行
UPDATE shop.orders SET status = 2 ... WHERE id = 1001后,1001最新一条消息的change_type=+U,status_name=SHIPPED。 - 插入
1003后,1003最新一条消息的change_type=+I,status_name=CREATED。 - 快照阶段
1002的status_name=PAID。
测试中还通过findLatestRecordByOrderId按order_id取同一业务键的最新记录来断言 CDC 增量语义,用convertHeadersToMap把 KafkaHeaders转为 Map 后校验 header 值,从实现层面印证了上述行为。
为什么这条链路这样写
1.Metadata只暴露了后续真的会用到的 CDC 字段
这份已验证配置里,只补了三个字段:source_database、source_table、change_type。
这不是巧合,而是有意为之:Metadatatransform 支持更丰富的元数据 Key,例如EventTime(数据变更事件时间戳)、Delay(采集延迟)、SourceTimestamp(源库提交时间戳)、以及 MySQL-CDC 专属的BinlogFile、BinlogPos、BinlogRow、Gtid(快照行为null)。本链路只投影后续 SQL 整形与 Kafka 消息头真正需要的三个字段,避免 payload 与 header 携带冗余信息。
2.Sql负责统一做业务整形
这次已验证链路里,SQL 实际产出的结果是:
- 快照阶段,
1001被写成status_name=CREATED - 快照阶段,
1002被写成status_name=PAID - 更新之后,
1001变成了status_name=SHIPPED - 插入之后,
1003被写成status_name=CREATED - 这次被断言的记录里还包含了
sync_source,其中快照阶段的1001还包含source_name
这说明CASE WHEN status = 0 THEN 'CREATED' ...的枚举翻译逻辑对快照行和 binlog 增量行均生效,CDC 行经过 Sql transform 后会携带一致的业务字段视图。
3.kafka_headers_fields会同时影响 header 和 payload
这个行为是在快照阶段1001那条记录上明确断言过的:source_database、source_table、change_type被写入 Kafka header 后,就不会继续保留在那条 JSON payload 里。也就是说,被列入kafka_headers_fields的字段会从消息 value 中"摘除",改由消息头承载。
从源码层面看,Kafka sink 在 KafkaSinkWriter.java 中还会对相关配置做前置校验:kafka_headers_fields不支持NATIVE格式(因为此时 key/value 已是byte[],headers 已编码在行内);partition_key_fields与kafka_headers_fields不能重叠;kafka_message_value_fields与kafka_headers_fields也不能重叠。设计这条链路时,把order_id留给分区键、把三个元数据字段留给消息头,恰好避开了这些冲突约束。
落地产出:从 Kafka 侧观察到的最终形态
综合上述配置与断言,recipe_mysql_orderstopic 中最终消息的形态可以概括为:
- 消息 key:由
partition_key_fields = ["order_id"]决定,用于分区路由与顺序保证。 - 消息 header:
source_database=shop、source_table=orders、change_type=+I/+U/-D(消费侧可通过 KafkaConsumerRecord.headers()读取,无需反序列化整个 JSON 即可路由)。 - 消息 value(JSON payload):只保留业务字段
order_id、order_no、user_id、amount、status_name,以及整形后的source_name、sync_source;三个元数据字段已不在 payload 中。
这种"轻量 payload + 富 header"的形态非常适合下游做基于消息头的过滤路由(如按来源库表分流、按变更类型触发不同处理),同时保持消息体只承载业务语义,是 CDC 数据入 Kafka 时一种值得复用的范式。
相关文档
- MySQL-CDC source 连接器:参数全表、MySQL 用户权限、binlog 开启方式、启动/停止模式、无主键表处理、server-id 冲突规避等。
- Kafka sink 连接器:
kafka_headers_fields、partition_key_fields、kafka_message_value_fields、format等参数的完整说明与示例。 - Metadata transform:全部元数据 Key(含 MySQL-CDC 专属的 Binlog/GTID 字段)与投影规则。
- Sql transform:SQL 表达式转换的语法与能力边界。
若要快速在自己的环境中复现,可直接复用 mysqlcdc_to_kafka_with_transforms.conf 这份作业配置,并参考 KafkaRecipeIT.java 中容器搭建、建表、CDC 用户授权与断言逻辑来准备环境。
- 数据集成
- ETL
- 大数据
- 批处理
- 流处理
- 变更数据捕获
【免费下载链接】seatunnel
SeaTunnel is a multimodal, high-performance, distributed, massive data integration tool.
相关推荐
SeaTunnel MySQL CDC 到 Kafka 实战:Metadata 元数据注入与 Kafka Headers 字段整形
SeaTunnel MySQL CDC 到 Kafka 实战:Metadata 元数据注入与 Kafka Headers 字段整形 本篇文章基于 SeaTunn
数据集成ETL大数据批处理流处理变更数据捕获SeaTunnel 实战:PostgreSQL CDC 实时同步到 Iceberg(字段整形与 Upsert 完整指南)
SeaTunnel 实战:PostgreSQL CDC 实时同步到 Iceberg(字段整形与 Upsert 完整指南) 本篇技术指南基于 Apache Sea
数据集成ETL大数据批处理流处理变更数据捕获SeaTunnel MySQL CDC 到 Elasticsearch 实战:过滤、转换与自定义字段的数据同步配方
SeaTunnel MySQL CDC 到 Elasticsearch 实战:过滤、转换与自定义字段的数据同步配方 导读 本文是 SeaTunnel 官方 Re
数据集成ETL大数据批处理流处理变更数据捕获
创作声明:本文部分内容由AI辅助生成(AIGC),仅供参考