ARTICLE DETAIL

资讯详情

深耕郑州网站建设与运营推广的一线实战洞察。

Canal数据同步实战:自定义Kafka消息格式与序列化器改造

Canal数据同步实战:自定义Kafka消息格式与序列化器改造 1. 项目缘起一次数据格式引发的“血案”最近在搞一个数据中台的项目需要把MySQL的变更数据实时同步到下游的十几个微服务里。技术选型上Canal监听MySQL的binlog然后投递到Kafka这几乎是业内的标准答案听起来很完美。我一开始也是这么想的直到下游的同事拿着Canal吐出来的JSON数据来找我眉头皱得能夹死苍蝇。“老哥你这数据格式我们没法直接用啊。”他指着屏幕说“你看data字段里是一个数组里面每个对象都带着完整的表结构字段但我们只需要id和update_timetype字段是INSERT、UPDATE、DELETE但我们希望是更业务化的CREATE、MODIFY、REMOVE还有这个es字段指executeTime我们想要的是标准的时间戳不是这个格式……”我一看确实。Canal默认的JSON格式是为了通用性设计的包含了数据库变更的完整元数据比如数据库名、表名、SQL类型、变更前/后的数据行等等。但对于具体的消费方来说他们往往只关心业务相关的核心字段并且希望格式符合自己系统的契约。这就好比厨房给你上了一整只没切分的烤鸭虽然原料顶级但你想直接卷饼吃还得自己动手片皮太麻烦了。这就是我们这次要解决的核心问题定制化Canal的输出。不是简单地用用就完事而是要深入其内部修改它投递到Kafka的消息体格式让它产出的“数据食粮”更符合下游各个“食客”的口味。这个过程涉及到对Canal客户端适配器、消息编码器乃至Kafka生产者配置的深度干预。下面我就把这次“庖丁解鸭”式的改造过程从原理到实操完整地拆解一遍。2. Canal与Kafka对接默认流程与核心痛点在动手改造之前我们必须先彻底理解Canal和Kafka在默认情况下是如何协同工作的。这就像医生动手术前必须清楚人体的解剖结构一样。2.1 默认数据流与JSON结构Canal的整体架构分为Server、Client和Adapter。我们通常说的“Canal”指的是Server它伪装成MySQL的Slave拉取binlog并解析成内部结构化的CanalEntry.Entry。而将Entry转化为具体目的地如Kafka消息的工作是由Client Adapter完成的。当你使用canal.adapter或canal.deployer中自带的canal-client时它会通过一个叫CanalKafkaProducer的类或类似实现来发送消息。默认情况下它使用一个SimpleMessageSerializer或类似的序列化器将CanalEntry.Entry转换成JSON字符串。一个典型的、未经处理的Canal-Kafka消息JSON格式如下{ data: [ { id: 1, name: test, create_time: 2023-10-27 12:00:00, update_time: 2023-10-27 12:00:00 } ], database: test_db, es: 1698393600000, id: 1, isDdl: false, mysqlType: { id: bigint(20), name: varchar(255), ... }, old: [ { name: old_test } ], pkNames: [id], sql: , sqlType: { id: -5, name: 12, ... }, table: user, ts: 1698393600123, type: UPDATE }字段解析与痛点data: 变更后的数据行列表。痛点永远是个数组即使单行操作包含全字段下游可能只需要其中几个。type: 操作类型固定为INSERT/UPDATE/DELETE。痛点无法自定义为业务术语。old: 仅UPDATE时存在表示被修改字段的旧值。痛点结构不一致有时是数组有时是对象下游解析麻烦。mysqlType/sqlType: 字段的MySQL和JDBC类型信息。痛点对绝大多数纯业务消费方无用徒增消息体积。es(executeTime): binlog中的执行时间戳。痛点格式可能是毫秒值也可能是其他格式下游需要统一。ts: Canal处理时间戳。痛点同上需要格式统一。2.2 为何默认格式常“不合身”这个默认格式的设计初衷是信息无损和通用性。它确保了任何下游系统无论其业务逻辑如何都能从这条消息中还原出一次完整的数据库变更事件。但这恰恰成了它在具体生产环境中的“阿喀琉斯之踵”网络与存储开销每条消息都携带了大量元数据mysqlType,sqlType,database,table等如果表字段很多单条消息体积可能膨胀数倍。在超大规模数据同步场景下这会给Kafka集群的带宽、磁盘以及下游消费者的反序列化性能带来不必要的压力。消费端解析复杂度下游业务程序员需要编写额外的代码来从data数组中提取所需字段判断old字段的存在性转换type枚举。这增加了业务代码的复杂度和出错概率。契约僵化默认格式是一个“霸王条款”所有消费者都必须接受。但当不同业务团队对数据有不同的格式要求例如用户服务需要user_id和email订单服务需要order_sn和amount时要么各自在消费端做转换重复劳动要么就需要我们在源头进行定制化分发。因此修改Canal的输出格式不是一个可有可无的优化而是在特定规模和数据使用场景下的必要架构决策。它的本质是在数据源头进行轻量的ETL提取、转换、加载实现“一发多收各取所需”的高效数据供给模式。3. 改造方案选型从“外敷”到“内服”的三种策略明确了问题接下来就是选择解决方案。根据对Canal架构的侵入程度和改造复杂度主要有三条路径我称之为“外敷”、“介入”和“内服”。3.1 方案一Kafka Connect 单消息转换SMT—— “外敷疗法”这是最“云原生”、对Canal最无侵入的方案。思路是Canal依然生产原始格式的消息到Kafka的一个原始主题如canal.raw.topic然后使用Kafka Connect框架搭配单消息转换Single Message Transform, SMT插件消费原始主题的消息进行格式转换再写入到另一个净化后的主题如canal.clean.topic供下游使用。优点完全解耦Canal和格式转换逻辑分离彼此独立部署、升级、扩缩容。灵活强大Kafka Connect生态丰富有现成的Cast,InsertField,ReplaceField,HoistField等SMT也可以通过编写自定义SMT实现复杂逻辑。可视化与管理一些平台如Confluent Platform提供了对Connect集群的可视化管理界面。缺点与实操考量架构复杂度引入了Kafka Connect集群这一新的中间件需要额外的运维成本。延迟增加数据流从Canal - Kafka(Topic A) - Connect - Kafka(Topic B) - Consumer比直接Canal - Kafka - Consumer多了一跳端到端延迟会增加几十到几百毫秒。资源消耗Connect集群本身需要消耗计算和内存资源。适用场景适合团队已有Kafka Connect技术栈或对Canal代码掌控力弱且可以接受额外延迟和复杂度的场景。对于追求极致实时性和架构简洁性的项目此方案需慎重。3.2 方案二定制Canal Client的MessageSerializer —— “介入疗法”这是最直接、最经典的改造方式。Canal Client在发送消息到Kafka前需要通过一个序列化器MessageSerializer将内部对象转为字节。我们可以实现一个自定义的序列化器在其中完成JSON格式的组装逻辑。核心步骤找到接口研究你使用的Canal Client版本如canal.client包找到MessageSerializer接口或类似接口如CanalMessageSerializer。实现类创建一个新类例如CustomCanalMessageSerializer实现该接口。在serializer方法中你拿到的是CanalEntry.Entry或CanalMessage对象这是最原始、信息最全的变更数据。定制组装在这个方法里你可以自由地从Entry中提取RowChange和RowData。只选取你需要的字段如id,name构建新的JSON对象。将INSERT/UPDATE/DELETE映射为CREATE/MODIFY/REMOVE。将时间戳格式化为yyyy-MM-dd HH:mm:ss或ISO8601字符串。过滤掉mysqlType、sqlType等无用信息。配置替换在Canal Client的配置文件中通常是application.yml或canal.properties将canal.mq.serializer或类似配置项的值从默认的org.apache.canal.client.impl.SimpleMessageSerializer改为你自定义类的全限定名。优点直击要害在数据产生的第一时间进行转换没有冗余流程。性能最优端到端路径最短延迟最低。掌控力强可以对Canal产生的原始数据结构进行任意操作。缺点与Canal版本绑定自定义序列化器依赖于Canal Client的内部API。如果Canal版本升级内部类结构或接口可能发生变化导致你的代码需要适配升级存在一定的维护成本。需要打包部署你需要将自定义的序列化器类打包进Jar并确保Canal服务能够加载到它通常放在lib目录或通过classpath指定。实操心得这是我最推荐大多数团队的方案。它平衡了效果、复杂度和可控性。在实现时务必在你的序列化器里做好异常捕获和日志记录因为这里一旦出错整条消息就会丢失。建议至少记录下出错的Entry的简要信息如tableName,eventType方便排查。3.3 方案三修改Canal Adapter源码并重编译 —— “内服疗法”这是最彻底、也是最“重”的方案。直接下载Canal的源码找到负责生成Kafka消息的模块通常是canal.adapter模块下的kafka相关代码直接修改其消息构建逻辑然后重新编译打包替换官方的发行版。优点终极定制你可以修改任何细节甚至改变整个处理流程。深度集成你的定制逻辑会成为Canal的一部分部署简单一个包。缺点维护噩梦你完全脱离了官方的主线版本。每次官方修复Bug或发布新特性你都需要手动合并代码冲突会非常多维护成本极高。技术门槛高需要深入理解Canal多个模块的代码结构。风险大自行修改可能引入未知的Bug且失去了官方社区的支持。结论除非你有非常特殊、稳定的定制需求且团队有强大的源码维护能力否则强烈不推荐此方案。这相当于维护一个自己的Canal分支代价巨大。综合来看方案二自定义MessageSerializer是性价比最高的选择。它既能实现深度定制又保持了与官方主线的可维护性关联。接下来我们就聚焦于方案二进行实战演练。4. 实战实现自定义MessageSerializer让我们一步步实现一个CustomKafkaMessageSerializer。假设我们的目标格式是{ operation: MODIFY, table: user, key: 1, change_time: 2023-10-27 12:00:00, after: { id: 1, username: new_name, status: 1 }, before: { username: old_name } }要求operation映射为业务术语只同步id, username, status字段时间格式化为字符串before只包含变更的字段。4.1 环境准备与依赖确认首先你需要一个可以编译Java项目的环境。确保你的Canal Client版本。这里以使用较广泛的canal.client为例。在你的项目pom.xml中引入对应版本的Canal Client依赖。例如对于1.1.7版本dependency groupIdcom.alibaba.otter/groupId artifactIdcanal.client/artifactId version1.1.7/version !-- 使用 provided 或 compile 范围取决于你如何部署 -- scopeprovided/scope /dependency !-- 还需要JSON处理库如Jackson -- dependency groupIdcom.fasterxml.jackson.core/groupId artifactIdjackson-databind/artifactId version2.15.0/version /dependency4.2 编写自定义序列化器创建一个类实现com.alibaba.otter.canal.client.kafka.MessageSerializer接口注意不同版本接口名或包名可能有差异请以实际代码为准。package com.yourcompany.canal.serializer; import com.alibaba.otter.canal.client.kafka.MessageSerializer; import com.alibaba.otter.canal.protocol.Message; import com.fasterxml.jackson.databind.ObjectMapper; import com.fasterxml.jackson.databind.node.ObjectNode; import org.apache.commons.lang.StringUtils; import org.slf4j.Logger; import org.slf4j.LoggerFactory; import java.text.SimpleDateFormat; import java.util.Date; import java.util.List; /** * 自定义Canal到Kafka的消息序列化器 */ public class CustomKafkaMessageSerializer implements MessageSerializer { private static final Logger LOGGER LoggerFactory.getLogger(CustomKafkaMessageSerializer.class); private static final ObjectMapper OBJECT_MAPPER new ObjectMapper(); private static final SimpleDateFormat DATE_FORMAT new SimpleDateFormat(yyyy-MM-dd HH:mm:ss); // 假设我们只关心这些表的这些字段 private static final String TABLE_USER user; private static final String[] USER_FIELDS {id, username, status}; Override public byte[] serialize(String destination, Message message) { if (message null || message.getId() -1) { return null; } try { Listcom.alibaba.otter.canal.protocol.FlatMessage flatMessages message.getFlatMessages(); if (flatMessages null || flatMessages.isEmpty()) { return null; } // 这里为了简化我们只处理第一条消息。实际生产环境可能需要遍历batch com.alibaba.otter.canal.protocol.FlatMessage flatMessage flatMessages.get(0); // 1. 构建根JSON对象 ObjectNode rootNode OBJECT_MAPPER.createObjectNode(); // 2. 映射操作类型 String opType mapOperationType(flatMessage.getType()); rootNode.put(operation, opType); // 3. 添加表名 rootNode.put(table, flatMessage.getTable()); // 4. 处理主键作为key (简化处理取第一个主键字段的第一个值) String key extractPrimaryKey(flatMessage); rootNode.put(key, key); // 5. 格式化时间戳 (使用Canal处理时间) rootNode.put(change_time, DATE_FORMAT.format(new Date(flatMessage.getTs()))); // 6. 构建after数据 (只取需要的字段) ObjectNode afterNode filterData(flatMessage.getData(), flatMessage.getTable()); if (afterNode ! null afterNode.size() 0) { rootNode.set(after, afterNode); } // 7. 构建before数据 (只取变更的字段) if (flatMessage.getOld() ! null !flatMessage.getOld().isEmpty()) { ObjectNode beforeNode filterOldData(flatMessage.getOld(), flatMessage.getTable()); if (beforeNode ! null beforeNode.size() 0) { rootNode.set(before, beforeNode); } } // 8. 转换为JSON字节 return OBJECT_MAPPER.writeValueAsBytes(rootNode); } catch (Exception e) { LOGGER.error(序列化Canal消息到自定义JSON格式失败, messageId: {}, destination: {}, message.getId(), destination, e); // 根据业务需求决定是抛出异常还是返回null或错误标记 // 抛出异常会导致整个batch发送失败返回null会忽略此条消息 return null; } } /** * 映射Canal操作类型到业务操作类型 */ private String mapOperationType(String canalType) { if (StringUtils.isEmpty(canalType)) { return UNKNOWN; } switch (canalType.toUpperCase()) { case INSERT: return CREATE; case UPDATE: return MODIFY; case DELETE: return REMOVE; default: return canalType; } } /** * 提取主键值 (简化版) */ private String extractPrimaryKey(com.alibaba.otter.canal.protocol.FlatMessage flatMessage) { if (flatMessage.getPkNames() ! null !flatMessage.getPkNames().isEmpty() flatMessage.getData() ! null !flatMessage.getData().isEmpty()) { String pkName flatMessage.getPkNames().get(0); Object firstRow flatMessage.getData().get(0); if (firstRow instanceof Map) { Object pkValue ((Map?, ?) firstRow).get(pkName); return pkValue ! null ? pkValue.toString() : ; } } return ; } /** * 过滤并构建新数据对象 */ private ObjectNode filterData(ListMapString, String dataList, String tableName) { if (dataList null || dataList.isEmpty()) { return null; } // 同样只处理第一行 MapString, String row dataList.get(0); ObjectNode node OBJECT_MAPPER.createObjectNode(); String[] fieldsToKeep getFieldsForTable(tableName); for (String field : fieldsToKeep) { if (row.containsKey(field)) { String value row.get(field); // 简单类型推断实际应根据mysqlType/sqlType处理 if (StringUtils.isNumeric(value)) { node.put(field, Long.parseLong(value)); } else { node.put(field, value); } } } return node; } /** * 过滤并构建旧数据对象 (只包含变更的字段) */ private ObjectNode filterOldData(ListMapString, String oldDataList, String tableName) { // 实现逻辑类似filterData但oldDataList的结构可能不同需要根据实际情况调整 // 这里是一个简化示例 if (oldDataList null || oldDataList.isEmpty()) { return null; } MapString, String oldRow oldDataList.get(0); ObjectNode node OBJECT_MAPPER.createObjectNode(); String[] fieldsToKeep getFieldsForTable(tableName); for (String field : fieldsToKeep) { if (oldRow.containsKey(field)) { node.put(field, oldRow.get(field)); } } return node; } private String[] getFieldsForTable(String tableName) { // 这里可以配置化从配置文件或数据库读取不同表需要同步的字段 if (TABLE_USER.equalsIgnoreCase(tableName)) { return USER_FIELDS; } // 默认返回空数组表示不同步任何字段 return new String[0]; } }代码关键点解析接口实现实现了MessageSerializer接口的serialize方法这是Kafka生产者调用的入口。数据源参数中的Message对象包含了FlatMessage列表这是Canal已经初步扁平化处理过的数据比原始Entry更易操作。类型安全与异常处理JSON构建过程被try-catch包裹任何异常都会记录日志并返回null导致该条消息被丢弃。在生产环境中这里需要更精细的错误处理策略比如将格式错误的消息投递到死信队列。字段过滤逻辑filterData和filterOldData方法根据表名决定保留哪些字段。这里写死了配置最佳实践是外部化配置例如从application.yml或Apollo配置中心读取。类型转换在filterData中我们做了一个简单的数字类型判断。实际上更准确的做法是结合FlatMessage中的mysqlType或sqlType字段进行精确的Java类型转换。4.3 配置与部署编写完代码并打包成Jar例如canal-custom-serializer-1.0.0.jar后需要让Canal服务加载它。步骤1放置Jar包将你的Jar包和它所依赖的第三方Jar包如Jackson放到Canal Server或Canal Adapter的lib目录下。例如/opt/canal-server/lib/。步骤2修改Canal配置编辑Canal Server的配置文件canal.properties或Adapter的application.yml找到Kafka生产者的序列化器配置项。对于Canal Servercanal.properties# 找到Kafka相关配置 canal.mq.servers kafka-broker1:9092,kafka-broker2:9092 canal.mq.topic your_topic # 关键配置指定自定义序列化器 canal.mq.serializer com.yourcompany.canal.serializer.CustomKafkaMessageSerializer # 确保使用flat message模式这样Message里才有FlatMessage列表 canal.mq.flatMessage true对于Canal Adapterapplication.ymlcanal.conf: mode: kafka mqServers: kafka-broker1:9092,kafka-broker2:9092 topic: your_topic # 关键配置 serializer: com.yourcompany.canal.serializer.CustomKafkaMessageSerializer flatMessage: true步骤3重启并验证重启Canal服务。然后对监听的MySQL表进行增删改操作使用Kafka控制台消费者或工具查看目标Topic的消息确认格式是否已按预期改变。踩坑记录我第一次部署时忘了把Jackson的依赖Jar包也放进lib目录导致Canal启动时报ClassNotFoundException。切记自定义序列化器及其所有非Canal内置的依赖都必须放入classpath通常是lib目录。可以使用maven-shade-plugin打成胖Jar或者手动管理所有依赖。5. 进阶动态配置与多Topic路由上面的示例是硬编码配置实际项目往往需要更灵活的策略不同表同步不同的字段甚至投递到不同的Kafka Topic。5.1 基于配置文件的动态规则我们可以创建一个配置文件如format-rules.yaml来定义规则rules: - table: ^test\.user$ # 正则匹配库名.表名 topic: topic_user fields: [id, username, email, status] operation_map: INSERT: USER_CREATED UPDATE: USER_UPDATED DELETE: USER_DELETED timestamp_field: update_time output_format: key: id wrap_object: true - table: ^test\.order$ topic: topic_order fields: [order_sn, user_id, amount, status] operation_map: INSERT: ORDER_CREATED UPDATE: ORDER_PAID # 不配置timestamp_field则使用Canal的ts output_format: key: order_sn wrap_object: false # 直接输出字段平铺的JSON然后在自定义序列化器的初始化阶段加载这个配置文件。在serialize方法中根据flatMessage.getDatabase()和flatMessage.getTable()匹配规则动态决定目标Topic甚至可以覆盖配置中的默认Topic。需要保留的字段列表。操作类型映射字典。输出JSON的结构。实现要点序列化器需要实现Configurable接口如果Canal支持或在初始化时从固定路径读取配置文件。同时要监听配置文件变化实现热更新。5.2 在序列化器中实现多Topic路由Canal默认一个实例或一个Adapter只能向一个固定Topic发送消息。要实现多Topic路由有两种思路思路A利用Kafka Producer的分区键Key和Topic前缀由下游消费者选择性消费。这并非真正的多Topic而是逻辑隔离。例如将所有消息发到canal_events主题但每条消息的Key设置为table_name。下游消费者可以使用Kafka的Consumer Group和分区分配策略或者自己过滤。这种方式简单但Topic内数据混杂不够清晰。思路B修改Canal Client使其支持在序列化器中动态返回目标Topic名。这需要更深入的改造。你需要研究Canal Client的Kafka生产者代码。通常MessageSerializer的serialize方法只负责生产消息体byte[]。Topic是在上层调用时决定的。你可能需要自定义一个MessageSerializer的子接口增加一个getTargetTopic(FlatMessage)方法。修改Canal Client中调用序列化器的代码先调用getTargetTopic获取Topic名再调用serialize获取消息体然后发送到对应的Topic。或者更“黑科技”一点在你的序列化器内部直接根据消息内容持有一个或多个KafkaProducer实例自己完成向不同Topic的发送。但这会严重破坏Canal原有的流程和事务语义风险极高不推荐。更优雅的方案如果多Topic路由是强需求可以考虑使用方案一Kafka Connect。让Canal先统一发到一个原始Topic然后在Connect中通过Router或自定义SMT根据消息内容将其路由到不同的目标Topic。这是Kafka生态更标准、更解耦的做法。6. 生产环境下的注意事项与优化将定制化的Canal投入生产还有一系列工程问题需要解决。6.1 监控与告警序列化错误率监控在你的CustomKafkaMessageSerializer中增加计数器统计序列化成功/失败的次数。通过JMX暴露指标或直接打印到日志由ELK收集并配置告警。一旦错误率超过阈值如0.1%立即告警。消息格式兼容性监控消费端在解析消息时如果遇到无法解析的格式如缺少必需字段也应记录日志并告警。这能及时发现序列化逻辑的Bug或配置错误。端到端延迟监控在消息体中加入一个源头时间戳如binlog的executeTime在消费端计算当前时间与它的差值监控数据同步的延迟。6.2 性能考量JSON库选型Jackson是性能非常好的选择。避免在序列化器中做复杂的字符串拼接务必使用ObjectMapper这样的专业库。对象复用ObjectMapper是线程安全的应该声明为static final复用。SimpleDateFormat是线程不安全的在并发环境下必须使用ThreadLocal包装或改用DateTimeFormatterJava 8。字段过滤开销如果字段过滤规则非常复杂如很多正则匹配可能会成为性能瓶颈。可以考虑在初始化时将规则编译成更高效的数据结构如Trie树或预编译的Pattern。6.3 兼容性与版本升级接口稳定性自定义序列化器强依赖Canal Client的API。在Canal升级时务必检查MessageSerializer、Message、FlatMessage等类是否有不兼容变更。配置回滚在更改序列化器配置或升级Jar包时做好回滚方案。可以先让新旧序列化器并行运行将消息同时发送到新老两个Topic验证无误后再切换消费者。6.4 消息大小与压缩自定义格式可能会改变消息大小。如果过滤了大量字段消息会变小如果增加了新的嵌套结构消息可能变大。需要关注Kafka Topic的压缩设置compression.type如gzip,snappy,lz4。对于文本格式的JSON启用压缩通常能获得不错的压缩比节省带宽和存储但会略微增加CPU开销。建议在测试环境对比开启压缩前后的吞吐量和CPU使用率找到平衡点。经过以上步骤一个高度定制化、贴合业务需求的Canal-Kafka数据通道就搭建完成了。从下游消费端的反馈来看他们不再需要编写冗长的数据清洗代码直接拿到了“开箱即用”的业务事件开发效率和数据链路可靠性都得到了显著提升。这个过程虽然需要一些前期的开发投入但对于一个长期运行、多团队协作的数据同步项目来说这份投入在维护阶段会带来持续的回报。
返回列表