- 数据工程
- 大数据
- 批处理
- 流处理
【免费下载链接】seatunnel
SeaTunnel is a next-generation super high-performance, distributed, massive data integration tool.
SeaTunnel 内置的canal_json格式是连接 MySQL 变更数据捕获(CDC)生态与数据集成管线的关键桥接层:它既能将 Canal 采集器生成的 MySQL Binlog 变更消息流(INSERT/UPDATE/DELETE)反序列化为 SeaTunnel 行数据,也能将 SeaTunnel 内部产生的变更事件重新编码为 Canal JSON 消息写回 Kafka。本文将带你掌握canal_json格式的全部配置项、底层解析原理与可落地的 Kafka 读写配置,并结合仓库源码与测试用例验证每一项行为。
什么是 Canal 格式
Canal(Changelog Data Capture)是一款能够将 MySQL 的数据变更以实时流的方式同步到其他系统的 CDC 工具。Canal 为 changelog(变更日志)提供了一套统一的格式规范,并支持使用JSON与protobuf两种方式序列化消息,其中 protobuf 是 Canal 的默认格式。
SeaTunnel 提供canal_json格式来实现两个方向的能力:
- 反序列化(Deserialization Schema):把 Canal 生成的 JSON 消息解释为 SeaTunnel 的 INSERT / UPDATE / DELETE 变更事件;
- 序列化(Serialization Schema):把 SeaTunnel 内部的 INSERT / UPDATE / DELETE 事件编码为 Canal JSON 消息,并投递到 Kafka 等下游存储。
需要特别注意的是一个已知限制:当前 SeaTunnel无法将 UPDATE_BEFORE 与 UPDATE_AFTER 合并为单条 UPDATE 消息。因此序列化时,SeaTunnel 会把 UPDATE_BEFORE 与 UPDATE_AFTER 分别编码为 Canal 的DELETE与INSERT消息(下文源码章节会展示这一映射逻辑)。
典型应用场景
将 Canal JSON 消息接入 SeaTunnel 后,可以在诸多场景中直接复用这套统一的变更日志语义:
- 将数据库的增量数据实时同步到其他系统(如数仓、消息队列、搜索引擎);
- 审计日志采集,记录每一行数据的变更轨迹;
- 基于数据库变更流构建实时物化视图;
- 对数据库表的历史变更进行 temporal join(时态关联)等。
格式选项(Format Options)
canal_json格式在 SeaTunnel 配置中通过format = canal_json启用,其配套选项定义在 CanalJsonFormatOptions.java 中,汇总如下:
| 选项 | 默认值 | 是否必填 | 说明 |
|---|---|---|---|
format | 无 | 是 | 指定数据格式,此处必须为canal_json |
canal_json.ignore-parse-errors | false | 否 | 解析出错时跳过出错字段与出错行而不是让作业失败;出错字段会被置为 null |
canal_json.database.include | 无 | 否 | 可选正则表达式,按 Canal 记录中的database元数据字段进行正则匹配,只读取特定数据库的 changelog 行;模式串与 Java 的Pattern兼容 |
canal_json.table.include | 无 | 否 | 可选正则表达式,按 Canal 记录中的table元数据字段进行正则匹配,只读取特定表的 changelog 行;模式串与 Java 的Pattern兼容 |
format
format是必填项,在 Kafka 连接器的MessageFormat枚举中对应CANAL_JSON取值(见 MessageFormat.java)。Kafka 连接器默认格式为json,因此显式声明format = canal_json是开启 Canal 语义解析的前提(见 Config.java 中format选项的定义)。
canal_json.ignore-parse-errors
该选项控制解析容错行为,默认false(解析失败直接抛错)。需要留意一个值得注意的源码细节:在 Kafka Source 的当前实现中,CANAL_JSON分支直接以setIgnoreParseErrors(true)构建反序列化器(见 KafkaSourceConfig.java),即 Kafka 消费侧默认容忍单条消息解析错误;而通用格式选项文档中的默认值仍是false。实际使用时应以你所使用的连接器实现与文档为准。
canal_json.database.include 与 canal_json.table.include
这两个选项通过正则表达式对 Canal 消息的database与table元数据字段做前置过滤。从 CanalJsonDeserializationSchema.java 可以看到,它们最终被编译为 JavaPattern,并在反序列化入口处逐条匹配:只有database与table均命中的消息才会继续向下游处理,未命中的直接跳过,从而实现"只同步指定库表"的精细化订阅。
反序列化原理:Canal 消息如何变成 SeaTunnel 变更事件
CanalJsonDeserializationSchema(源码)是反序列化的核心实现,其处理流程如下:
- 若配置了
database/table正则,先对消息中的database、table元数据字段做匹配过滤,不匹配直接返回; - 读取
data与type两个关键字段,按type分发处理; - 对每条数据行调用 JSON 反序列化器转换为
SeaTunnelRow,并设置对应的RowKind与tableId。
事件类型与 RowKind 的映射
SeaTunnel 内部用RowKind(+I表示 INSERT,-U/+U表示 UPDATE_BEFORE / UPDATE_AFTER,-D表示 DELETE)来表达变更语义。反序列化时按 Canal 的type字段做如下转换:
Canaltype | SeaTunnel 输出 | 说明 |
|---|---|---|
INSERT | 一条+I行 | 将data数组中每条记录直接收集输出 |
UPDATE | 一对-U(UPDATE_BEFORE)与+U(UPDATE_AFTER)行 | 从data解析变更后值、从old解析变更前值,成对输出 |
DELETE | 一条-D行 | 将data中记录标记为删除 |
CREATE/ALTER/QUERY | 跳过 | 这类 DDL 或查询事件data为 null,直接忽略 |
一个非常实用的实现细节体现在 UPDATE 处理上:Canal 的old数组里只包含被修改的字段,未修改的字段并不会出现。SeaTunnel 在生成 UPDATE_BEFORE 行时,会检查old中缺失的字段,并把 UPDATE_AFTER 行中对应字段的值回填进 UPDATE_BEFORE,保证"变更前"行是一份完整记录(见 CanalJsonDeserializationSchema.java)。
DDL 与非数据事件的处理
Canal 消息中的isDdl: true、data: null的 CREATE / ALTER / QUERY 类事件会被直接跳过,不会产生数据行;但如果一条 INSERT / UPDATE / DELETE 消息的data字段为空,则会抛出IllegalStateException(提示 "Null data value ... Cannot send downstream"),因为这类事件无法向下游产出有效数据。
错误处理与容错
当解析过程中抛出运行时异常时:
ignoreParseErrors = false(默认):抛出统一的jsonOperationError异常,作业失败,便于及时感知问题;ignoreParseErrors = true:吞掉该条消息的解析错误,跳过出错字段/行继续处理,出错字段置为 null。
序列化原理:SeaTunnel 变更事件如何编码为 Canal JSON
CanalJsonSerializationSchema(源码)负责把 SeaTunnel 行数据编码为 Canal JSON。它输出的是一个极简的两字段结构:{"data": {...}, "type": "INSERT"}——即数据主体放在data下,操作类型放在type下。
关键的 RowKind 映射逻辑如下(见rowKind2String方法):
| SeaTunnel RowKind | 编码后的 Canaltype |
|---|---|
| INSERT | INSERT |
| UPDATE_AFTER | INSERT |
| UPDATE_BEFORE | DELETE |
| DELETE | DELETE |
这正是文档开头所述限制的落地实现:由于无法合并 UPDATE_BEFORE / UPDATE_AFTER,编码时干脆把两者分别降级为 DELETE 与 INSERT。从源码注释可以看到,序列化时也主动丢弃了database、ts、old等 Canal 元数据字段,只保留data与type两个最小必要字段。
实战:通过 Kafka 读写 Canal JSON 消息
以下场景来自 SeaTunnel 官方文档示例:假设 Canal 正在捕获 MySQLinventory库中products表(包含id、name、description、weight四列)的变更,并将变更事件写入 Kafka topicproducts_binlog。SeaTunnel 消费该 topic、把变更事件解释为行数据后,再以canal_json格式写回另一个 Kafka topic。
一条完整的 UPDATE 消息长什么样
下面这条消息是一次 UPDATE 变更事件:products表中id = 111的那一行,weight字段值从5.15被更新为5.18:
{ "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" }该消息中各字段含义如下(各字段的完整语义可参考 Canal 官方文档):data为变更后的行数据数组;old为变更前的被修改字段(old中只包含发生变化的字段);type为操作类型;database/table为变更来源的库表名;es/ts为事件时间戳;isDdl标识是否为 DDL 事件;mysqlType/sqlType为列类型信息;pkNames为主键字段列表;sql为 DDL 原始 SQL。SeaTunnel 反序列化时主要消费data、old、type、database、table五个字段,其余字段由 Canal 侧携带但不参与下游数据构建。
SeaTunnel 作业配置示例
假设上述消息已同步到 Kafka topicproducts_binlog,可以用下面的 SeaTunnel 配置消费该 topic 并解释变更事件,同时把处理结果以canal_json格式写入下游 Kafka topicconsume-binlog:
env { parallelism = 1 job.mode = "BATCH" } source { Kafka { bootstrap.servers = "kafkaCluster:9092" topic = "products_binlog" result_table_name = "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 会据此构建CanalJsonDeserializationSchema(见 KafkaSourceConfig.java);schema.fields声明的字段名与类型需要与data数组中的 JSON 对象字段一一对应,解析时按此 schema 完成类型转换;start_mode = earliest表示从最早位点开始消费,便于演示。 - Sink 侧:
format = canal_json会让 Kafka Sink 使用CanalJsonSerializationSchema(见 DefaultSeaTunnelRowSerializer.java),将行数据编码为{"data": {...}, "type": "..."}结构投递到consume-binlog。
序列化输出的实际效果
以仓库测试资源 canal-data-filter-table.txt 中的真实数据为例,反序列化后再序列化得到的 Canal JSON 形如:
{"data":{"id":106,"name":"hammer","description":"18oz carpenter hammer","weight":1.0},"type":"INSERT"} {"data":{"id":106,"name":"hammer","description":null,"weight":1.0},"type":"DELETE"}可以看到,weight从 1.0 改为 1.0 的 UPDATE 事件(原始 Canal 消息为type: UPDATE+old字段),经 SeaTunnel 序列化后被拆成了先DELETE后INSERT两条消息,与文档所述的 UPDATE_BEFORE / UPDATE_AFTER 降级策略完全一致。
源码与测试验证
canal_json格式的实现与验证集中在seatunnel-formats/seatunnel-format-json模块:
- 格式选项定义:CanalJsonFormatOptions.java —— 定义了
database.include、table.include、ignore-parse-errors三个选项的 key 与描述; - 反序列化实现:CanalJsonDeserializationSchema.java —— 事件类型分发、库表正则过滤、
old字段回填、错误处理等核心逻辑; - 序列化实现:CanalJsonSerializationSchema.java —— RowKind 到
INSERT/DELETE的映射; - 单元测试:CanalJsonSerDeSchemaTest.java —— 覆盖了库表正则过滤(
testFilteringTables,使用^my.*与^prod.*模式)、空消息、非法 JSON、data缺失、未知操作类型等异常路径,并对 27 条真实 Canal 消息完成了反序列化→序列化的往返断言; - 测试数据:canal-data-filter-table.txt —— 包含 INSERT / UPDATE / DELETE / CREATE 等多种类型的真实 Canal JSON 消息。
测试中的预期输出(如kind=+I、kind=-U、kind=+U、kind=-D)直接印证了本文所述的 RowKind 映射关系:INSERT 产生+I,UPDATE 产生-U与+U成对输出,DELETE 产生-D。
注意事项与已知限制
- UPDATE 消息无法合并:SeaTunnel 序列化时无法把 UPDATE_BEFORE 与 UPDATE_AFTER 合并为一条 Canal UPDATE,只能输出 DELETE + INSERT 两条消息。若下游依赖 Canal 原生的 UPDATE 语义,需要在应用层自行处理;
- Schema 必须预先声明:消费 Canal JSON 消息时,必须在 source 的
schema.fields中声明与data对应的字段及类型,SeaTunnel 不会自动推导 Canal 消息中的列结构; - 库表过滤是前置匹配:
canal_json.database.include/canal_json.table.include使用 Java 正则完整匹配(Pattern.matcher(...).matches()),未命中的消息会被整条跳过,可用于降低无效数据的传输与解析开销; - DDL 事件不产出数据:
isDdl: true的 CREATE / ALTER 事件会被忽略,SeaTunnel 只消费数据变更事件; - 解析容错因连接器而异:格式选项文档中的默认容错为关闭(
false),但 Kafka Source 当前实现默认以ignoreParseErrors = true构建 canal_json 反序列化器,排查解析异常时需要同时考虑连接器层的这一行为。
- 数据工程
- 大数据
- 批处理
- 流处理
【免费下载链接】seatunnel
SeaTunnel is a next-generation super high-performance, distributed, massive data integration tool.
相关推荐
SeaTunnel Canal JSON 格式解析与实战:基于 Canal CDC 消息的 MySQL 增量同步指南
SeaTunnel Canal JSON 格式解析与实战:基于 Canal CDC 消息的 MySQL 增量同步指南 Canal 是阿里开源的 CDC(Chan
数据集成ETL大数据批处理流处理变更数据捕获Flink Canal Format 实战指南:基于 canal-json 的 MySQL CDC 变更数据捕获与同步
Flink Canal Format 实战指南:基于 canal json 的 MySQL CDC 变更数据捕获与同步 Canal 是阿里巴巴开源的 CDC(C
后端大数据流处理批处理告别JSON解析难题:Canal完美适配MySQL 8.0 JSON格式binlog全指南
告别JSON解析难题:Canal完美适配MySQL 8.0 JSON格式binlog全指南 你是否还在为MySQL 8.0中JSON字段的binlog解析而头疼
后端变更数据捕获数据同步数据集成
创作声明:本文部分内容由AI辅助生成(AIGC),仅供参考