SeaTunnel Canal JSON 格式解析与实战:基于 Canal CDC 消息的 MySQL 增量同步指南
【免费下载链接】seatunnelSeaTunnel is a multimodal, high-performance, distributed, massive data integration tool.项目地址: https://gitcode.com/GitHub_Trending/se/seatunnel
Canal 是阿里开源的 CDC(Change Data Capture,变更数据捕获)工具,能够实时将 MySQL 的 binlog 变更流式同步到其他系统,并对外输出统一的 changelog(变更日志)格式。SeaTunnel 通过canal_json格式在 Source/Sink 与 Canal 消息之间建立桥接:既能把 Canal 输出的 JSON 消息反序列化为 SeaTunnel 内部的 INSERT/UPDATE/DELETE 行数据,也能把 SeaTunnel 的行变更重新编码为 Canal JSON 消息写入 Kafka 等存储。读完本文,你将掌握canal_json格式的完整配置参数、Canal 消息字段语义、底层反序列化/序列化实现原理,以及一套可直接运行的 Kafka 消费与投递实战配置。
Canal JSON 格式是什么
Canal JSON 是 Canal 组件为 changelog 提供的一种统一格式 Schema。Canal 原生支持使用 JSON 和 protobuf(默认)两种方式序列化消息,SeaTunnel 的canal_json格式专门处理其中的 JSON 形态。
SeaTunnel 会把 Canal JSON 消息解释为 SeaTunnel 行模型中的三种变更类型(对应 RowKind 语义):
- INSERT:新增行;
- UPDATE:更新行(SeaTunnel 内部拆分为 UPDATE_BEFORE 与 UPDATE_AFTER 两行);
- DELETE:删除行。
这种能力在以下场景中非常实用:
- 将数据库增量数据同步到其他系统(如异构数据库、数仓、消息队列);
- 审计日志采集与分析;
- 基于数据库变更构建实时物化视图;
- 对数据库表的变更历史做时态关联(temporal join)等。
反向地,SeaTunnel 也支持把内部的行变更编码为 Canal JSON 消息,投递到 Kafka 等存储。需要注意:当前 SeaTunnel 无法把 UPDATE_BEFORE 和 UPDATE_AFTER 合并为单条 UPDATE 消息,因此在编码时会将 UPDATE_BEFORE 和 UPDATE_AFTER 分别编码为 DELETE 与 INSERT 两条 Canal 消息(当开启合并选项时可输出 UPDATE,详见下文序列化原理)。
格式选项(Format Options)
canal_json格式在 SeaTunnel 配置中通过format = canal_json启用,支持的选项如下:
| Option | 默认值 | 是否必填 | 说明 |
|---|---|---|---|
format | (无) | 是 | 指定使用的格式,此处必须为canal_json |
canal_json.ignore-parse-errors | false | 否 | 解析出错时跳过对应字段和行而不是使任务失败;出错字段会被置为 null |
canal_json.database.include | (无) | 否 | 可选正则表达式,仅读取指定数据库的 changelog 行,通过匹配 Canal 记录中的database元字段实现;模式串与 Java 的Pattern兼容 |
canal_json.table.include | (无) | 否 | 可选正则表达式,仅读取指定表的 changelog 行,通过匹配 Canal 记录中的table元字段实现;模式串与 Java 的Pattern兼容 |
以上四个选项在源码中有直接对应实现:CanalJsonFormatOptions.java 中定义了database.include、table.include两个字符串选项,并复用了通用 JSON 格式的IGNORE_PARSE_ERRORS布尔选项(默认false)。其中canal_json.ignore-parse-errors虽然默认关闭,但在处理脏数据较多的生产 binlog 流时,开启它可以避免单条异常消息导致整个作业失败。
Canal JSON 消息结构详解
Canal 输出的 changelog 消息是一个结构化的 JSON 对象。以下是一条从 MySQLproducts表捕获到的 UPDATE 操作消息(该表有id、name、description、weight四列):
{ "data": [ { "id": "111", "name": "scooter", "description": "Big 2-wheel scooter", "weight": "5.18" } ], "database": "inventory", "es": 1589373560000, "id": 9, "isDdl": false, "mysqlType": { "id": "INTEGER", "name": "VARCHAR(255)", "description": "VARCHAR(512)", "weight": "FLOAT" }, "old": [ { "weight": "5.15" } ], "pkNames": [ "id" ], "sql": "", "sqlType": { "id": 4, "name": 12, "description": 12, "weight": 7 }, "table": "products", "ts": 1589373560798, "type": "UPDATE" }这条消息的含义是:products表中id = 111的行的weight字段值从5.15变更为5.18。各关键字段语义如下:
| 字段 | 含义 |
|---|---|
data | 变更后(after)的数据数组,每行一个对象;字段名对应表列名,值为字符串 |
old | 变更前(before)的数据数组,仅包含发生变化的字段;UPDATE/DELETE 消息中携带 |
database | 变更所属的数据库名,供canal_json.database.include过滤使用 |
table | 变更所属的表名,供canal_json.table.include过滤使用 |
type | 变更类型,取值为INSERT/UPDATE/DELETE等 |
ts | 事件时间戳(毫秒),SeaTunnel 会将其写入行的 event time 元数据 |
es、id、isDdl、mysqlType、pkNames、sql、sqlType | Canal 附带的其他元信息(执行时间、消息 ID、是否 DDL、MySQL 类型、主键列、SQL、JDBC 类型),SeaTunnel 反序列化时主要关注data、old、type、database、table、ts六个字段 |
反序列化实现原理(Deserialization)
SeaTunnel 对 Canal JSON 的反序列化由 CanalJsonDeserializationSchema.java 实现,其核心处理流程可以从源码中梳理出来:
- 元数据过滤:若配置了
database/table正则,先对消息的database、table字段做Pattern.matcher(...).matches()全量匹配,不匹配的消息直接丢弃(对应canal_json.database.include/canal_json.table.include)。 - DDL 事件跳过:当
data字段为 null 时,如果操作类型是QUERY、CREATE、ALTER(DDL 类事件),直接跳过;否则抛出Null data value ... Cannot send downstream异常。 - 按操作类型分发:
INSERT:将data数组中的每一行转换为SeaTunnelRow输出;UPDATE:同时解析data(after)与old(before)数组,为 before 行补齐未变化字段(old中不存在的字段从 after 行拷贝),然后分别以RowKind.UPDATE_BEFORE和RowKind.UPDATE_AFTER输出两行;DELETE:将data中的每一行标记为RowKind.DELETE输出;- 其他未知操作类型抛出
Unknown operation type异常。
- 时间戳注入:若消息携带
ts字段,通过MetadataUtil.setEventTime(row, ts)将事件时间写入行的元数据,供后续窗口、时态关联等使用。 - 错误处理:整个解析过程包裹在 try-catch 中,当
ignoreParseErrors = false时抛出jsonOperationError使任务失败;为true时静默跳过异常消息。
SeaTunnel 使用 CanalJsonSerDeSchemaTest.java 对上述行为进行了完整验证,测试覆盖了表过滤(testFilteringTables)、空 data 行、非 JSON 输入、空 JSON、无 data 字段、未知操作类型等边界场景,以及多行 DELETE 事件批量的反序列化结果,可作为理解语义的补充参考。
序列化实现原理(Serialization)
SeaTunnel 将内部行变更编码为 Canal JSON 消息由 CanalJsonSerializationSchema.java 实现。输出的 Canal 消息包含old、data、type、database、table、ts六个字段。
类型映射的核心逻辑在rowKind2String方法中:
INSERT→type: "INSERT";UPDATE_AFTER→ 默认编码为type: "INSERT"(即原文所述:UPDATE 被拆成 DELETE + INSERT 两条消息);UPDATE_BEFORE、DELETE→type: "DELETE"。
同时源码支持一个mergeUpdateEventFlag合并开关:当开启时,收到UPDATE_BEFORE会先缓存在cacheUpdateBeforeRow中并暂不输出,待收到紧随其后的UPDATE_AFTER时,将 before 行放入old数组、after 行放入data数组,输出一条type: "UPDATE"的完整消息。这为需要下游精确 UPDATE 语义的场景提供了底层支持。
database和table字段来源于SeaTunnelRow的 tableId(通过TablePath解析出库名与表名),ts字段来源于行的 event time 选项。Kafka 连接器侧通过 DefaultSeaTunnelRowSerializer.java 调用该序列化器,实现"读 Canal JSON → 写 Canal JSON"的消息格式透传。
实战:Kafka 消费与投递示例
假设 Canal 已把 MySQLproducts表的变更同步到 Kafka 的products_binlogtopic,我们可以用下面的 SeaTunnel 配置消费该 topic 并解释变更事件,再把结果投递到另一个 Kafka topic(consume-binlog),形成一条完整的 binlog 搬运链路:
env { parallelism = 1 job.mode = "BATCH" } source { Kafka { bootstrap.servers = "kafkaCluster:9092" topic = "products_binlog" plugin_output = "kafka_name" start_mode = earliest schema = { fields { id = "int" name = "string" description = "string" weight = "string" } }, format = canal_json } } transform { } sink { Kafka { bootstrap.servers = "localhost:9092" topic = "consume-binlog" format = canal_json } }配置要点说明:
- Source 端:
format = canal_json让 Kafka Source 把 topic 中的 Canal JSON 消息反序列化为 SeaTunnel 行;schema.fields声明了products表的物理列(id、name、description、weight),SeaTunnel 会从data数组中按字段名映射取值,类型声明为int/string等 SeaTunnel 类型。Kafka Source 对canal_json格式的装配逻辑可在 KafkaSourceConfig.java 中找到(其引入了CanalJsonDeserializationSchema)。 - Sink 端:
format = canal_json将行变更重新编码为 Canal JSON 消息写入下游 topic。若希望下游得到真正的UPDATE语义(而不是拆分的 DELETE + INSERT),可结合上述mergeUpdateEventFlag的合并机制。 - BATCH 模式:示例使用
job.mode = "BATCH",如需持续消费 binlog 流,应改为STREAMING。
更多关于 SeaTunnel 内部行模型与外部数据编码(Avro、Debezium JSON、Protobuf 等)的映射关系,可参考 Formats 总览 与 Data Format Handling;该文档原文位于 canal-json.md。
常见问题与注意事项
- UPDATE 语义的拆分问题:SeaTunnel 默认将 UPDATE_BEFORE / UPDATE_AFTER 编码为 DELETE / INSERT 两条 Canal 消息,依赖 Canal 原生 UPDATE 语义的下游系统需要对此兼容,或利用合并开关输出完整 UPDATE。
- 类型映射:Canal 消息中
data的字段值以字符串形式出现(如"5.18"),需要在schema.fields中声明目标类型,SeaTunnel 会按声明类型转换;未声明或类型不符的字段在开启canal_json.ignore-parse-errors时会被置为 null。 - 过滤粒度:
database.include/table.include使用 Java 正则的matches()全匹配语义,配置时需写全完整匹配表达式(例如products只匹配库/表名恰好为products的情况,需要前缀通配时写.*products.*)。 - DDL 与查询事件:
QUERY/CREATE/ALTER等非数据变更事件会被 SeaTunnel 静默跳过,不会阻塞下游管道。
【免费下载链接】seatunnelSeaTunnel is a multimodal, high-performance, distributed, massive data integration tool.项目地址: https://gitcode.com/GitHub_Trending/se/seatunnel
创作声明:本文部分内容由AI辅助生成(AIGC),仅供参考