尧图网站设计 尧图网站设计YAOTU DESIGN
ARTICLE DETAIL

资讯详情

深耕网站设计与一线实操的经验洞察。

SeaTunnel Canal JSON 格式解析与实战:基于 Canal CDC 消息的 MySQL 增量同步指南

SeaTunnel Canal JSON 格式解析与实战:基于 Canal CDC 消息的 MySQL 增量同步指南 SeaTunnel Canal JSON 格式解析与实战基于 Canal CDC 消息的 MySQL 增量同步指南【免费下载链接】seatunnelSeaTunnel is a multimodal, high-performance, distributed, massive data integration tool.项目地址: https://gitcode.com/GitHub_Trending/se/seatunnelCanal 是阿里开源的 CDCChange 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 Optionscanal_json格式在 SeaTunnel 配置中通过format canal_json启用支持的选项如下Option默认值是否必填说明format无是指定使用的格式此处必须为canal_jsoncanal_json.ignore-parse-errorsfalse否解析出错时跳过对应字段和行而不是使任务失败出错字段会被置为 nullcanal_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、sqlTypeCanal 附带的其他元信息执行时间、消息 ID、是否 DDL、MySQL 类型、主键列、SQL、JDBC 类型SeaTunnel 反序列化时主要关注data、old、type、database、table、ts六个字段反序列化实现原理DeserializationSeaTunnel 对 Canal JSON 的反序列化由 CanalJsonDeserializationSchema.java 实现其核心处理流程可以从源码中梳理出来元数据过滤若配置了database/table正则先对消息的database、table字段做Pattern.matcher(...).matches()全量匹配不匹配的消息直接丢弃对应canal_json.database.include/canal_json.table.include。DDL 事件跳过当data字段为 null 时如果操作类型是QUERY、CREATE、ALTERDDL 类事件直接跳过否则抛出Null data value ... Cannot send downstream异常。按操作类型分发INSERT将data数组中的每一行转换为SeaTunnelRow输出UPDATE同时解析dataafter与oldbefore数组为 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 事件批量的反序列化结果可作为理解语义的补充参考。序列化实现原理SerializationSeaTunnel 将内部行变更编码为 Canal JSON 消息由 CanalJsonSerializationSchema.java 实现。输出的 Canal 消息包含old、data、type、database、table、ts六个字段。类型映射的核心逻辑在rowKind2String方法中INSERT→type: INSERTUPDATE_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 topicconsume-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、weightSeaTunnel 会从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),仅供参考
返回列表