资讯中心

SeaTunnel Canal JSON 格式完全指南:基于 MySQL Binlog 的 CDC 数据读写实战

📅 2026/9/27 21:12:28
SeaTunnel Canal JSON 格式完全指南:基于 MySQL Binlog 的 CDC 数据读写实战
数据工程大数据批处理流处理【免费下载链接】seatunnelSeaTunnel is a next-generation super high-performance, distributed, massive data integration tool.项目地址https://gitcode.com/gh_mirrors/sea/seatunnel点击查看免费下载SeaTunnel 内置的canal_json格式是连接 MySQL 变更数据捕获CDC生态与数据集成管线的关键桥接层它既能将 Canal 采集器生成的 MySQL Binlog 变更消息流INSERT/UPDATE/DELETE反序列化为 SeaTunnel 行数据也能将 SeaTunnel 内部产生的变更事件重新编码为 Canal JSON 消息写回 Kafka。本文将带你掌握canal_json格式的全部配置项、底层解析原理与可落地的 Kafka 读写配置并结合仓库源码与测试用例验证每一项行为。什么是 Canal 格式CanalChangelog 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 Optionscanal_json格式在 SeaTunnel 配置中通过format canal_json启用其配套选项定义在 CanalJsonFormatOptions.java 中汇总如下选项默认值是否必填说明format无是指定数据格式此处必须为canal_jsoncanal_json.ignore-parse-errorsfalse否解析出错时跳过出错字段与出错行而不是让作业失败出错字段会被置为 nullcanal_json.database.include无否可选正则表达式按 Canal 记录中的database元数据字段进行正则匹配只读取特定数据库的 changelog 行模式串与 Java 的Pattern兼容canal_json.table.include无否可选正则表达式按 Canal 记录中的table元数据字段进行正则匹配只读取特定表的 changelog 行模式串与 Java 的Pattern兼容formatformat是必填项在 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 内部用RowKindI表示 INSERT-U/U表示 UPDATE_BEFORE / UPDATE_AFTER-D表示 DELETE来表达变更语义。反序列化时按 Canal 的type字段做如下转换CanaltypeSeaTunnel 输出说明INSERT一条I行将data数组中每条记录直接收集输出UPDATE一对-UUPDATE_BEFORE与UUPDATE_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 JSONCanalJsonSerializationSchema源码负责把 SeaTunnel 行数据编码为 Canal JSON。它输出的是一个极简的两字段结构{data: {...}, type: INSERT}——即数据主体放在data下操作类型放在type下。关键的 RowKind 映射逻辑如下见rowKind2String方法SeaTunnel RowKind编码后的 CanaltypeINSERTINSERTUPDATE_AFTERINSERTUPDATE_BEFOREDELETEDELETEDELETE这正是文档开头所述限制的落地实现由于无法合并 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-binlogenv { 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.javaschema.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: UPDATEold字段经 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 消息。测试中的预期输出如kindI、kind-U、kindU、kind-D直接印证了本文所述的 RowKind 映射关系INSERT 产生IUPDATE 产生-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 反序列化器排查解析异常时需要同时考虑连接器层的这一行为。赞分享数据工程大数据批处理流处理【免费下载链接】seatunnelSeaTunnel is a next-generation super high-performance, distributed, massive data integration tool.项目地址https://gitcode.com/gh_mirrors/sea/seatunnel点击查看免费下载相关推荐SeaTunnel Canal JSON 格式解析与实战基于 Canal CDC 消息的 MySQL 增量同步指南SeaTunnel Canal JSON 格式解析与实战基于 Canal CDC 消息的 MySQL 增量同步指南 Canal 是阿里开源的 CDCChan数据集成ETL大数据批处理流处理变更数据捕获Flink Canal Format 实战指南基于 canal-json 的 MySQL CDC 变更数据捕获与同步Flink Canal Format 实战指南基于 canal json 的 MySQL CDC 变更数据捕获与同步 Canal 是阿里巴巴开源的 CDCC后端大数据流处理批处理告别JSON解析难题Canal完美适配MySQL 8.0 JSON格式binlog全指南告别JSON解析难题Canal完美适配MySQL 8.0 JSON格式binlog全指南 你是否还在为MySQL 8.0中JSON字段的binlog解析而头疼后端变更数据捕获数据同步数据集成上一篇终极指南如何在Windows 10/11上完美运行Android应用下一篇Privacy Badger工作原理深度解析从启发式算法到数据结构的完整揭秘创作声明:本文部分内容由AI辅助生成(AIGC),仅供参考

看完文章,想为自己的企业也做一次专业网站诊断?

尧图顾问免费为您评估现有网站,并给出建站/改版建议与报价方案。

免费获取方案