news 2026/9/27 21:12:25

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

作者头像

张小明

前端开发工程师

1.2k 24
文章封面图
SeaTunnel Canal JSON 格式完全指南:基于 MySQL Binlog 的 CDC 数据读写实战
  • 数据工程
  • 大数据
  • 批处理
  • 流处理

【免费下载链接】seatunnel

SeaTunnel 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 格式

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-errorsfalse否解析出错时跳过出错字段与出错行而不是让作业失败;出错字段会被置为 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(源码)是反序列化的核心实现,其处理流程如下:

  1. 若配置了database/table正则,先对消息中的database、table元数据字段做匹配过滤,不匹配直接返回;
  2. 读取data与type两个关键字段,按type分发处理;
  3. 对每条数据行调用 JSON 反序列化器转换为SeaTunnelRow,并设置对应的RowKind与tableId。

事件类型与 RowKind 的映射

SeaTunnel 内部用RowKind(+I表示 INSERT,-U/+U表示 UPDATE_BEFORE / UPDATE_AFTER,-D表示 DELETE)来表达变更语义。反序列化时按 Canal 的type字段做如下转换:

CanaltypeSeaTunnel 输出说明
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
INSERTINSERT
UPDATE_AFTERINSERT
UPDATE_BEFOREDELETE
DELETEDELETE

这正是文档开头所述限制的落地实现:由于无法合并 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。

注意事项与已知限制

  1. UPDATE 消息无法合并:SeaTunnel 序列化时无法把 UPDATE_BEFORE 与 UPDATE_AFTER 合并为一条 Canal UPDATE,只能输出 DELETE + INSERT 两条消息。若下游依赖 Canal 原生的 UPDATE 语义,需要在应用层自行处理;
  2. Schema 必须预先声明:消费 Canal JSON 消息时,必须在 source 的schema.fields中声明与data对应的字段及类型,SeaTunnel 不会自动推导 Canal 消息中的列结构;
  3. 库表过滤是前置匹配:canal_json.database.include/canal_json.table.include使用 Java 正则完整匹配(Pattern.matcher(...).matches()),未命中的消息会被整条跳过,可用于降低无效数据的传输与解析开销;
  4. DDL 事件不产出数据:isDdl: true的 CREATE / ALTER 事件会被忽略,SeaTunnel 只消费数据变更事件;
  5. 解析容错因连接器而异:格式选项文档中的默认容错为关闭(false),但 Kafka Source 当前实现默认以ignoreParseErrors = true构建 canal_json 反序列化器,排查解析异常时需要同时考虑连接器层的这一行为。
  • 数据工程
  • 大数据
  • 批处理
  • 流处理

【免费下载链接】seatunnel

SeaTunnel is a next-generation super high-performance, distributed, massive data integration tool.

项目地址:https://gitcode.com/gh_mirrors/sea/seatunnel
点击查看免费下载

相关推荐

上一篇:终极指南:如何在Windows 10/11上完美运行Android应用
下一篇:Privacy Badger工作原理深度解析:从启发式算法到数据结构的完整揭秘

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

版权声明: 本文来自互联网用户投稿,该文观点仅代表作者本人,不代表本站立场。本站仅提供信息存储空间服务,不拥有所有权,不承担相关法律责任。如若内容造成侵权/违法违规/事实不符,请联系邮箱:809451989@qq.com进行投诉反馈,一经查实,立即删除!
网站建设 2026/9/27 21:11:57

不会代码?3步搞定公司网站建设推广,免费工具全攻略

不会代码?3步搞定公司网站建设推广,免费工具全攻略 自己不会写代码,但老板让你下周把公司官网弄上线,还得兼顾SEO和域名备案?别慌,这事儿真没那么难。过去搞网站得找外包花几万块,现在用对 免费工具 ,普通人也能在三天内把 公司网站建设推广 的全流程跑通。…

作者头像 李华
网站建设 2026/9/27 21:11:44

海南免税店网上商城怎么选服务商,避开拖一周改需求的大坑

海南免税店网上商城怎么选服务商,避开拖一周改需求的大坑 改个首页Banner图,建站公司拖你整整一周? 这种经历在海南免税店网上商城开发圈太常见了。 很多老板想给线下免税店开个线上入口,结果选错了外包团队。 改需求像挤牙膏,上线前发现漏洞满屏飞。 今天不聊虚的,直接拆解海南免税店网上商城 怎么选…

作者头像 李华
网站建设 2026/9/27 21:11:34

图解 React 源码系列:react-illustration-series 原理学习路线与源码导读

教程前端 【免费下载链接】react-illustration-series 图解react源码, 用大量配图的方式, 致力于将react原理表述清楚. 项目地址: https://gitcode.com/gh_mirrors/re/react-illustration-series 点击查看 免费下载 本文以 react-illustration-series 的系列总览文…

作者头像 李华
网站建设 2026/9/27 21:11:21

太原seo推广优化避坑指南3步搞定域名服务器速查手册

太原seo推广优化避坑指南3步搞定域名服务器速查手册 域名买错,服务器配错,网站上线三天没流量?别急着骂人,先看看你的后台配置。 很多太原的老板找我们做网站,第一句话不是问价格,而是问:“我那个域名怎么解析不生效?”或者“服务器为什么老连不上?”这背后其实是 域名服务器搞不懂…

作者头像 李华
网站建设 2026/9/27 21:11:07

2026最新企业网站开发的公司怎么选?防黑挂马实战指南

2026最新企业网站开发的公司怎么选?防黑挂马实战指南 网站被黑挂马,后台密码全失效,页面跳出赌博广告,这时候你慌不慌?很多老板第一反应是删库重装,但这往往治标不治本。2026年最新的网络攻击手段已经进化到利用AI批量生成漏洞利用代码,传统防火墙几乎失效。作为在华东地区深耕十年的建站从业者,我见过太…

作者头像 李华
网站建设 2026/9/27 21:10:51

天津建站模板源码避坑指南新手入门实测

天津建站模板源码避坑指南新手入门实测 找天津建站公司,最怕的就是被忽悠花大钱买个“半成品”。很多新手入门第一反应是搜“天津建站模板源码”,结果发现报价从几千到几万不等,功能看着差不多,价格差出十倍。 别急着付钱。在天津本地做网站,尤其是企业站或商城,…

作者头像 李华