首页/新闻资讯/正文详情

从 MySQL CDC 到 Debezium-JSON:SeaTunnel compatible_debezium_json 格式的使用与原理全解析

发布时间:2026/9/27 9:36:54 来源:云帆数科 栏目:资讯中心
从 MySQL CDC 到 Debezium-JSON:SeaTunnel compatible_debezium_json 格式的使用与原理全解析
数据工程大数据批处理流处理【免费下载链接】seatunnelSeaTunnel is a next-generation super high-performance, distributed, massive data integration tool.项目地址https://gitcode.com/gh_mirrors/sea/seatunnel点击查看免费下载SeaTunnel 提供了一种名为compatible_debezium_json的数据格式用于将 CDCChange Data Capture连接器内部产生的 Debezium 变更记录重新序列化为 Debezium-JSON 消息并发布到 Kafka 等消息队列系统。本文以 官方格式文档 为主线结合仓库内格式模块、CDC 连接器与 Kafka 连接器的源码实现完整讲解该格式的配置方法、消息结构、底层转换原理与验证方式帮助你快速搭建“MySQL CDC → Kafka”的 Debezium 兼容数据链路。一、为什么需要 compatible_debezium_jsonSeaTunnel 的 CDC 连接器如 MySQL-CDC内部基于 Debezium 捕获数据库变更默认会把这些变更解析为 SeaTunnel 自身的SeaTunnelRow行数据供下游 Transform 或 Sink 消费。但在很多场景下下游系统希望拿到的是标准 Debezium-JSON 消息下游已经围绕 Debezium 生态构建了消费链路例如基于 Kafka Connect 的消息格式约定、Debezium UI、Canal/OGG 之外的消息治理工具需要把捕获到的原始变更事件含before/after、操作类型op、变更来源source、时间戳等元数据原样透传给消息中间件由其他团队按 Debezium 规范解析需要在一套消息格式标准下接入多个来源库保持消息结构与 Debezium 原生输出一致。compatible_debezium_json正是为此设计它把 CDC 源内部产生的SourceRecordDebezium 标准事件载体用 Kafka Connect 的JsonConverter序列化成 JSON 字符串并以topic、key、value三个字段组成的行数据向下游输出从而与 Debezium 生态天然兼容。二、快速上手MySQL CDC → Kafka 的完整配置原始文档给出了一个开箱即用的配置示例下面完整保留并逐段解释。env { parallelism 1 job.mode STREAMING checkpoint.interval 15000 } source { MySQL-CDC { result_table_name table1 base-urljdbc:mysql://localhost:3306/test startup.modeINITIAL table-names[ database1.t1, database1.t2, database2.t1 ] # compatible_debezium_json options format compatible_debezium_json debezium { # include schema into kafka message key.converter.schemas.enable false value.converter.schemas.enable false # include ddl include.schema.changes true # topic prefix database.server.name mysql_cdc_1 } } } sink { Kafka { source_table_name table1 bootstrap.servers localhost:9092 # compatible_debezium_json options format compatible_debezium_json } }env 段parallelism 1将作业并行度设为 1保证变更事件按 binlog 顺序写出CDC 场景下通常需要保序job.mode STREAMING作业以流式模式运行持续消费增量变更checkpoint.interval 15000每 15 秒做一次 checkpoint用于故障恢复与进度管理。source 段MySQL-CDCbase-urlJDBC 连接地址指向要监听的 MySQL 实例startup.mode INITIAL启动时先做全量快照snapshot再进入增量监听阶段table-names要监听的表列表支持跨库多表如database1.t1、database1.t2、database2.t1format compatible_debezium_json核心开关把 CDC 源内部记录输出格式切换为 Debezium-JSON 兼容格式debezium {...}以debezium为前缀的 Debezium 客户端属性透传块见下一节详解。sink 段Kafkasource_table_name table1消费上游 source 注册的结果表bootstrap.serversKafka 集群地址format compatible_debezium_json让 Kafka Sink 以“上游行中携带的topic/key/value直接作为 Kafka 消息的 topic/key/value 写出”的方式序列化。三、配置参数详解source 侧format在 CDC 连接器公共配置中format定义于 SourceOptions.java类型枚举DeserializeFormat定义于 DeserializeFormat.java取值default默认输出 SeaTunnel 行数据或compatible_debezium_json效果format的枚举标识compatible_debezium_json与格式模块中的标识符一一对应CompatibleDebeziumJsonDeserializationSchema.IDENTIFIER compatible_debezium_json。当设置为compatible_debezium_json后CDC 源不再按表的物理结构产出SeaTunnelRow而是统一产出一个三字段行topic、key、value全部为字符串类型。这一点在 IncrementalSource.java 的getProducedCatalogTables()中有明确实现当format为COMPATIBLE_DEBEZIUM_JSON时返回的 CatalogTable 使用CompatibleDebeziumJsonDeserializationSchema.DEBEZIUM_DATA_ROW_TYPE即topic、key、value三个字符串字段。source 侧debezium属性透传块debezium子配置块在 SourceOptions.java 中定义为DEBEZIUM_PROPERTIESmapType其语义是“以debezium为前缀的 Debezium 客户端属性”。示例中出现的几个属性作用如下属性示例值作用key.converter.schemas.enablefalse是否在 Kafka 消息的key中携带 Schema 结构Kafka Connect JsonConverter 的schemas.enable。设为false输出纯 JSON不携带schema信封value.converter.schemas.enablefalse是否在 Kafka 消息的value中携带 Schema 结构。设为false时 value 为纯 Debezium envelope JSONinclude.schema.changestrue是否将 DDLSchema 变更事件作为独立消息发出便于下游感知表结构变更database.server.namemysql_cdc_1Debezium 逻辑实例名是消息 topic 命名的前缀例如mysql_cdc_1.database1.t1同时用于区分不同 CDC 实例需要注意的是key.converter.schemas.enable与value.converter.schemas.enable两个开关的默认值是true见 DebeziumJsonDeserializeSchema.java 中的getOrDefault(..., true)即不配置时会在 JSON 中携带完整的 Schema 结构示例中显式设为false输出更精简的纯 JSON 消息。sink 侧formatKafka Sink 的format选项定义于 Config.java默认值为json支持compatible_debezium_json等格式。当 Kafka Sink 配置为compatible_debezium_json时序列化逻辑与默认 JSON 格式完全不同详见下一节源码剖析。四、源码级原理剖析4.1 格式模块的整体结构该格式独立成模块位于 seatunnel-formats/seatunnel-format-compatible-debezium-json包含三个核心类类角色CompatibleDebeziumJsonDeserializationSchema.java反序列化侧把 CDC 源产生的 DebeziumSourceRecord转成(topic, key, value)三字段的SeaTunnelRowCompatibleDebeziumJsonSerializationSchema.java序列化侧把(topic, key, value)行中指定的key或value字段原样输出为字节DebeziumJsonConverter.java核心转换器基于 Kafka ConnectJsonConverter完成SourceRecord的 key/value JSON 序列化4.2 反序列化侧SourceRecord → SeaTunnelRowCompatibleDebeziumJsonDeserializationSchema定义了输出行结构public static final String FIELD_TOPIC topic; public static final String FIELD_KEY key; public static final String FIELD_VALUE value; public static final SeaTunnelRowType DEBEZIUM_DATA_ROW_TYPE new SeaTunnelRowType( new String[] {FIELD_TOPIC, FIELD_KEY, FIELD_VALUE}, new SeaTunnelDataType[] { BasicType.STRING_TYPE, BasicType.STRING_TYPE, BasicType.STRING_TYPE });其deserialize(SourceRecord record)方法将每条 Debezium 变更事件转换为一行数据String key debeziumJsonConverter.serializeKey(record); String value debeziumJsonConverter.serializeValue(record); Object[] fields new Object[] {record.topic(), key, value}; return new SeaTunnelRow(fields);即行内topic字段来自SourceRecord.topic()key与value由转换器序列化产生。注意该实现仅面向SourceRecordCDC 内部事件面向字节流的deserialize(byte[])直接抛出UnsupportedEncodingException说明该格式只在 CDC 链路内部使用不用于通用的字节消息解析。4.3 转换器基于 Kafka Connect JsonConverterDebeziumJsonConverter是理解本格式的关键。它使用 Kafka Connect 的JsonConverter并依据keySchemaEnable/valueSchemaEnable两个开关选择内部转换方法开关为true时调用convertToJsonWithEnvelope输出带schema和payload信封的 JSON开关为false时调用convertToJsonWithoutEnvelope输出纯 JSON。同时无论开关如何转换器都强制设置DecimalFormat.NUMERIC使DECIMAL类型以数字而非字符串形式输出。例如测试用例 TestDebeziumJsonConverter.java 验证了BigDecimal.valueOf(1101, 2)即 11.01序列化后输出{k:11.01}、{v:11.01}而不是{k:11.01}。此外转换器对key 为 null的场景做了保护当SourceRecord没有 key例如主键为空的表或未配置 key 提取时serializeKey返回null避免空指针异常对应测试 testDebeziumSerializeKeyIsNull。4.4 CDC 连接器侧的集成方式以 MySQL-CDC 为例MySqlIncrementalSource.java 的createDebeziumDeserializationSchema()中if (DeserializeFormat.COMPATIBLE_DEBEZIUM_JSON.equals( config.get(JdbcSourceOptions.FORMAT))) { return (DebeziumDeserializationSchemaT) new DebeziumJsonDeserializeSchema( config.get(JdbcSourceOptions.DEBEZIUM_PROPERTIES)); }即当format compatible_debezium_json时使用DebeziumJsonDeserializeSchema而非默认的SeaTunnelRowDebeziumDeserializeSchema。DebeziumJsonDeserializeSchema见 源码从debezium配置块中读取key.converter.schemas.enable与value.converter.schemas.enable默认true构造CompatibleDebeziumJsonDeserializationSchema并把每个SourceRecord转成(topic, key, value)行收集输出。从源码结构可以推断CDC 连接器家族MySQL、PostgreSQL、Oracle、SQL Server、MongoDB 等都基于这套公共DebeziumJsonDeserializeSchema机制因此该格式对多个 CDC 数据源是通用的。4.5 Kafka Sink 侧的集成方式Kafka Sink 的序列化器 DefaultSeaTunnelRowSerializer.java 在format compatible_debezium_json时做了三项特殊处理topic 提取不再使用配置中固定的 topic而是直接取行内topic字段的值作为 Kafka 目标 topic即 CDC 源产生的SourceRecord.topic()如mysql_cdc_1.database1.t1key 序列化使用CompatibleDebeziumJsonSerializationSchema(rowType, true)输出行内key字段的字节内容value 序列化使用CompatibleDebeziumJsonSerializationSchema(rowType, false)输出行内value字段的字节内容。CompatibleDebeziumJsonSerializationSchema的实现非常直接构造时通过rowType.indexOf(isKey ? FIELD_KEY : FIELD_VALUE)定位目标列serialize()时取出该列的字符串并转成字节。由此可见从 CDC 源到 Kafka Sinkcompatible_debezium_json是一条“原样透传”链路CDC 内部把 Debezium 事件 JSON 化并填入(topic, key, value)行Kafka Sink 再把这三个字段直接写回 Kafka 消息的三要素保证消息与 Debezium 原生输出一致。五、Kafka 中的实际消息形态以key.converter.schemas.enable false、value.converter.schemas.enable false为例落到 Kafka 的每条消息大致形态如下具体字段以 Debezium 连接器实际输出为准示例仅供理解消息 key主键 JSON{id:1}消息 valueDebezium envelope含变更前后值与操作类型{ before: null, after: {id: 1, name: seatunnel, amount: 11.01}, source: {version: 1.9.7.Final, connector: mysql, db: database1, table: t1, ts_ms: 1718000000000}, op: c, ts_ms: 1718000000000 }其中before/after变更前 / 变更后的行数据null表示该侧无数据INSERT 的 before 为 nullDELETE 的 after 为 nullsource变更来源元数据连接器版本、库表名、binlog 位点时间戳等字段内容由 Debezium 连接器决定op操作类型ccreate、uupdate、ddelete、rread快照阶段读取。从源码看CDC 公共反序列化逻辑 SeaTunnelRowDebeziumDeserializeSchema.java 正是依据Envelope.OperationCREATE/READ/DELETE/UPDATE来区分这些操作并提取Envelope.FieldName.BEFORE/AFTER对应的结构体ts_ms变更时间戳毫秒。若将key/value.converter.schemas.enable设为true默认值则 key/value 会包裹在{schema: {...}, payload: {...}}信封结构中携带完整的 Connect Schema 描述便于严格类型的下游如 Schema Registry 类消费端解析。两种形态各有利弊纯 JSON 更精简易读带 Schema 的形态信息更完整实际使用中按下游消费能力取舍。若开启include.schema.changes trueDDL 事件如ALTER TABLE也会作为一类特殊消息发出下游可以据此感知并同步表结构变更。六、验证与测试格式模块自带单元测试 TestDebeziumJsonConverter.java覆盖两个关键行为Decimal 以数字输出构造包含Decimal.builder(2)字段的SourceRecord设置key/value为BigDecimal.valueOf(1101, 2)断言serializeKey输出{k:11.01}、serializeValue输出{v:11.01}验证DecimalFormat.NUMERIC生效key 为 null 的容错构造只有 value 没有 key 的SourceRecord断言serializeKey返回null而serializeValue正常输出{v:DebeziumTest}验证空 key 场景不抛异常。此外仓库 seatunnel-e2e 下的连接器端到端测试体系如 Kafka、CDC 相关 e2e 模块可结合真实 Kafka 实例验证完整链路的消息产出适合作为本配置的集成验证手段。七、适用场景与注意事项适用场景需要让下游直接消费 Debezium 标准 JSON 消息复用 Debezium 生态的消费端、治理端工具多源库统一以 Debezium-JSON 规范接入 Kafka消息结构一致便于统一解析与归档需要把 CDC 原始事件含 before/after、op、source 元数据完整透传而不做字段裁剪或结构转换。注意事项topic 由上游决定format compatible_debezium_json时Kafka Sink 的目标 topic 来自行内topic字段即database.server.name前缀 库表名Kafka Sink 配置中的topic参数不再生效需要注意 topic 的命名与权限规划key 可能为 null对于没有主键或未定义 key 提取规则的表消息 key 可能为空下游消费端需容忍空 keySchema 开关影响消息体积key/value.converter.schemas.enable默认true会显著增大消息体积如仅需纯 JSON请在debezium块中显式置为false即文档示例的做法并行度与保序CDC 场景建议保持较低并行度如示例中的parallelism 1以维持 binlog 顺序DDL 消息开启include.schema.changes true后DDL 事件也会进入消息流下游解析时需对这类消息做分支处理该格式面向 CDC 内部链路反序列化实现只接受SourceRecord不适用于将 Kafka 中的字节消息直接解析成该格式请勿在普通 Source 上使用。八、延伸阅读原始格式文档docs/en/connector-v2/formats/cdc-compatible-debezium-json.md格式模块源码seatunnel-formats/seatunnel-format-compatible-debezium-jsonCDC 公共格式选项SourceOptions.javaCDC 反序列化格式枚举DeserializeFormat.javaKafka Sink 序列化实现DefaultSeaTunnelRowSerializer.javaKafka Sink 配置文档docs/en/connector-v2/sink/Kafka.md赞分享数据工程大数据批处理流处理【免费下载链接】seatunnelSeaTunnel is a next-generation super high-performance, distributed, massive data integration tool.项目地址https://gitcode.com/gh_mirrors/sea/seatunnel点击查看免费下载相关推荐SeaTunnel CDC 兼容 Debezium-JSON 格式将 MySQL CDC 变更记录发布到 KafkaSeaTunnel CDC 兼容 Debezium JSON 格式将 MySQL CDC 变更记录发布到 Kafka 导读 SeaTunnel 的 compa数据集成ETL大数据批处理流处理变更数据捕获Flink CDC MySQL Binlog解析从Row格式到事件转换Flink CDC MySQL Binlog解析从Row格式到事件转换 一、MySQL Binlog与CDC技术基础 MySQL Binlog二进制日志是后端数据集成大数据流处理变更数据捕获数据同步SeaTunnel Debezium JSON 格式解析与生成数据库 CDC 变更事件SeaTunnel Debezium JSON 格式解析与生成数据库 CDC 变更事件 本篇聚焦 SeaTunnel 的 Debezium JSON 格式数据集成ETL大数据批处理流处理变更数据捕获上一篇detector.py全解析TACO模型训练、测试与评估完整流程下一篇如何用asynq实现高效请求批处理从入门到精通的完整教程创作声明:本文部分内容由AI辅助生成(AIGC),仅供参考

相关推荐

Google Antigravity SDK联网搜索实战:内置Web工具让AI Agent实时获取网络信息
Google Antigravity SDK联网搜索实战:内置Web工具让AI Agent实时获取网络信息

Google Antigravity SDK联网搜索实战:内置Web工具让AI Agent实时获取网络信息 【免费下载链接】antigravity-sdk-python A Python library for building AI agents that leverage the full power of Google Antigravity. 项目地址: https://gitcode.com/gh_mirror… · 2026/9/27 9:36:48

[人工智能]Python04:Numpy.linalg 实战指南
[人工智能]Python04:Numpy.linalg 实战指南

numpy.linalg 实战指南一份实用 numpy.linalg 指南,覆盖向量、矩阵乘法、线性系统、分解、特征问题、最小二乘和数值稳定性。理解 NumPy 线性代数背后的几何和形状约定。根据问题选择 solve、lstsq、QR、SVD 或特征值方法。检查残差、秩、条件数和数值容差。编写高效… · 2026/9/27 9:36:48

# 鸿蒙端侧 AI:把大模型装在设备本地,开启万物智联新范式 >
# 鸿蒙端侧 AI:把大模型装在设备本地,开启万物智联新范式 >

摘要 传统 AI 大多依赖云端服务器运算,一旦断网就失效,同时隐私数据上传存在风险。鸿蒙的核心思路是AI 能力下沉到端侧,将轻量化盘古大模型深度集成进操作系统底座,依托麒麟 NPU、MindSpore Lite 推理引擎,在手机、平板… · 2026/9/27 9:36:42

烧录程序版本管理实战:防止芯片烧错固件的关键策略
烧录程序版本管理实战:防止芯片烧错固件的关键策略

烧录程序版本管理这件事,我做了十几年嵌入式,踩过的坑比很多新手写过的代码都多。程序写错了改一改就行,芯片烧错了版本,轻则功能不符重则整板报废,而且往往是在批量产线上出问题,一发现就是几百片返工。我… · 2026/9/27 10:23:55

KubeVela vNext 多租户设计解析:KEP-2.14 中的 TenantDefinition 与 Tenant 平台原语
KubeVela vNext 多租户设计解析:KEP-2.14 中的 TenantDefinition 与 Tenant 平台原语

云原生DevOps运维微服务 【免费下载链接】kubevela The Modern Application Platform. 项目地址: https://gitcode.com/gh_mirrors/ku/kubevela 点击查看 免费下载 KEP-2.14 是 KubeVela vNext Roadmap(design/vela-core/keps/README.md)中的… · 2026/9/27 10:23:49

嵌入式驱动从“能跑”到“不崩”:量产级工程化实战指南
嵌入式驱动从“能跑”到“不崩”:量产级工程化实战指南

我见过太多这样的场景:嵌入式驱动开发出来的demo驱动在试验板上跑得稳稳的,功能正常、数据正确,开发人员自信满满地交出去,结果样机一到客户手里,要么几天崩一次,要么高低温环境一跑就死机,要么… · 2026/9/27 10:23:49

预处理、编译、汇编、链接
预处理、编译、汇编、链接

1.翻译环境和运行环境在 ANSI C 的任何⼀种实现中,存在两个不同的环境。第一种是翻译环境,在这个环境中,源代码被转换为可执行的机器指令(二进制指令);第二种是执行环境,它用于实际执行代码。1.… · 2026/9/27 10:23:49

KubeVela KEP-2.4 解读:Dispatcher 可插拔交付机制与跨集群分发设计
KubeVela KEP-2.4 解读:Dispatcher 可插拔交付机制与跨集群分发设计

云原生DevOps运维微服务 【免费下载链接】kubevela The Modern Application Platform. 项目地址: https://gitcode.com/gh_mirrors/ku/kubevela 点击查看 免费下载 本文基于 design/vela-core/keps/2.4-dispatchers/README.md 整理。该 KEP 当前标注为 Drafting&am… · 2026/9/27 10:23:43

Thonny在树莓派上的窗口机制与性能调优指南
Thonny在树莓派上的窗口机制与性能调优指南

1. 为什么树莓派用户总在Thonny里卡住——不是IDE太慢,是窗口逻辑没吃透你刚把树莓派4B插上电,烧好Raspberry Pi OS,双击桌面那个蓝色图标打开Thonny,结果发现:代码能写,但“运行”按钮点了没反应&#xff… · 2026/9/27 10:23:43

MATLAB雷达信号脉冲压缩仿真:LFM线性调频、匹配滤波与距离分辨率实现
MATLAB雷达信号脉冲压缩仿真:LFM线性调频、匹配滤波与距离分辨率实现

简介:这套Matlab仿真工具完整呈现雷达信号脉冲压缩过程,从线性调频(LFM)信号生成、目标回波仿真到匹配滤波压缩处理均有可运行代码支撑,面向电子信息工程、计算机、数学等专业学生,适用于课程设计、期末大作… · 2026/9/27 0:00:01

汕头网站建设制作厂家避坑指南:5大注意事项救急
汕头网站建设制作厂家避坑指南:5大注意事项救急

汕头网站建设制作厂家避坑指南:5大注意事项救急 改个需求建站公司拖一周,这种憋屈事我见得太多了。 很多汕头老板找本地建站团队,签合同前看着方案挺美,一上线就变脸。 今天不聊虚的,直接拆解找 汕头网站建设制作厂家 时的5个核心 注意事项… · 2026/9/27 0:00:01

多模态虚假新闻检测实战:BERT+ResNet双塔与对比学习
多模态虚假新闻检测实战:BERT+ResNet双塔与对比学习

简介:基于PyTorch的多模态虚假新闻检测项目完整代码包,面向自然语言处理与计算机视觉交叉方向的开发者、科研人员及毕业设计选题者,解决社交媒体中文本与图像联合识别虚假新闻的问题。系统以BERT预训练模型提取文本语义特征,以Res… · 2026/9/27 0:00:01

MATLAB雷达信号脉冲压缩仿真:LFM线性调频、匹配滤波与距离分辨率实现
MATLAB雷达信号脉冲压缩仿真:LFM线性调频、匹配滤波与距离分辨率实现

简介:这套Matlab仿真工具完整呈现雷达信号脉冲压缩过程,从线性调频(LFM)信号生成、目标回波仿真到匹配滤波压缩处理均有可运行代码支撑,面向电子信息工程、计算机、数学等专业学生,适用于课程设计、期末大作… · 2026/9/27 0:00:01

汕头网站建设制作厂家避坑指南:5大注意事项救急
汕头网站建设制作厂家避坑指南:5大注意事项救急

汕头网站建设制作厂家避坑指南:5大注意事项救急 改个需求建站公司拖一周,这种憋屈事我见得太多了。 很多汕头老板找本地建站团队,签合同前看着方案挺美,一上线就变脸。 今天不聊虚的,直接拆解找 汕头网站建设制作厂家 时的5个核心 注意事项… · 2026/9/27 0:00:01

多模态虚假新闻检测实战:BERT+ResNet双塔与对比学习
多模态虚假新闻检测实战:BERT+ResNet双塔与对比学习

简介:基于PyTorch的多模态虚假新闻检测项目完整代码包,面向自然语言处理与计算机视觉交叉方向的开发者、科研人员及毕业设计选题者,解决社交媒体中文本与图像联合识别虚假新闻的问题。系统以BERT预训练模型提取文本语义特征,以Res… · 2026/9/27 0:00:01

了解更多?预约专属演示

我们的顾问将为您一对一讲解产品与方案

企业微信二维码