Canal数据同步实战:定制Kafka JSON格式适配下游消费
1. 项目概述从Canal到Kafka的JSON格式改造最近在搞一个数据同步的项目核心是把MySQL的变更数据实时推送到下游系统。阿里开源的Canal是个老牌工具了它伪装成MySQL的Slave解析binlog然后把变更事件推出来用起来确实方便。我们一开始的架构很简单Canal Server解析数据通过内置的Kafka Producer把数据以默认的JSON格式扔到Kafka里下游的消费者再去消费处理。但做着做着就发现这个默认的JSON格式有点“水土不服”。下游的消费方不止一个有做实时数仓的有做缓存更新的还有做业务通知的。它们对数据格式的要求五花八门有的希望字段名是蛇形命名snake_case有的希望时间戳是毫秒级的整型还有的只关心某几个特定字段希望消息体尽可能精简。Canal默认输出的那个JSON结构是固定的包含了库名、表名、事件类型、变更前后的数据全集等等信息很全但也很“重”每次传输和解析都有不小的开销。所以这个项目的核心目标就明确了改造Canal输出到Kafka的消息格式。不是简单地换换字段名而是要设计一套灵活、可配置的JSON序列化方案让Canal吐出来的数据能更好地适配下游多样化的消费需求。这背后涉及到对Canal内部消息结构的理解、对Kafka生产者序列化器的定制以及对不同业务场景下数据契约的权衡。接下来我就把这次改造过程中的思路、踩过的坑和最终的方案详细拆解一遍。2. 核心需求与方案选型背后的逻辑为什么要改JSON格式这得从数据对接的痛点说起。Canal默认的格式我们可以称之为“原始Binlog镜像格式”。它长这样{ data: [ { id: 1, name: test, create_time: 2023-10-01 12:00:00 } ], database: test_db, table: user, type: INSERT, es: 1696147200000, ts: 1696147200123, sql: , old: null }这个格式的问题很具体字段命名风格不统一数据库字段是create_time但有些下游Java服务更习惯createTime这就导致了反序列化时的额外映射配置。时间格式可读性强但处理麻烦create_time字段是字符串下游每次用都要做解析。对于高吞吐场景这个解析开销累积起来不容忽视。信息冗余database,table,sql这些字段对于某些只需要纯数据内容的消费者来说是多余的增加了网络传输和存储的成本。结构僵化data是一个数组即使单条变更也是数组。有些消费者希望直接拿到平铺的JSON对象而不是多套一层。基于这些痛点我们的需求可以归纳为三点可定制化字段名、格式、轻量化减少不必要字段、结构化友好适配不同消费端的数据结构要求。方案选型上主要有三个方向方案一在下游消费者端做转换。这是最省事的办法不动Canal让每个消费者自己去解析原始格式然后转换成自己需要的格式。优点是架构改动最小。缺点也非常明显转换逻辑分散每个消费者都要写一遍重复劳动而且如果原始格式变更所有消费者都要跟着改维护成本高。方案二使用Canal的Adapter或ETL工具。Canal官方提供了Adapter组件支持将数据同步到MySQL、RocketMQ等也支持简单的表映射。社区也有基于Flink、Spark做ETL的方案。这个方案的优点是功能强大可以进行复杂的转换和清洗。但缺点是引入了新的、更重的系统组件增加了架构的复杂度和运维成本。对于我们这种格式转换需求明确但不算特别复杂的场景有点“杀鸡用牛刀”。方案三定制Canal的Kafka消息序列化器。这是最直接、侵入性最小、也最符合“单一职责”原则的方案。Canal在向Kafka发送消息时需要将内部的CanalEntry.Entry对象序列化成字节数组。默认使用的是StringSerializer它会把一个固定的JSON字符串发出去。我们可以实现一个自定义的Serializer在这个环节按照我们定义的规则将CanalEntry.Entry转换成我们想要的JSON格式。我们最终选择了方案三。理由很清晰它把格式转换的逻辑收口在了数据出口Canal Server下游消费者拿到就是“熟数据”开箱即用。整个架构没有增加新的移动部件只是替换了Kafka生产者的一小部分逻辑风险可控改造范围明确。接下来我们就深入这个自定义序列化器的实现细节。3. 深入Canal消息结构与序列化入口要定制输出首先得彻底搞清楚Canal内部的数据结构。Canal的核心数据单元是CanalEntry.Entry它是一个Protocol Buffers定义的消息。我们通过Kafka Producer发送的主要就是这种Entry。一个Entry包含了一次数据变更的完整上下文其关键结构如下Header: 包含日志文件名、位置等元信息。EntryType: 事件类型如TRANSACTIONBEGIN, ROWDATA, TRANSACTIONEND等。我们最关心的是ROWDATA。StoreValue: 实际存储的值。对于ROWDATA类型这里存放的是序列化后的RowChange消息。RowChange消息是重中之重它里面包含了tableId: 表ID。eventType: 具体的DML操作类型如INSERT、UPDATE、DELETE。rowDatas: 一个列表包含变更前后的行数据。每个RowData包含beforeColumns和afterColumns两个列表分别代表变更前和变更后的列值。每个Column包含index索引、sqlTypeJDBC类型、name列名、isKey是否主键、updated是否更新、isNull、value字符串格式的值。默认的JSON序列化器比如SimpleMessageSerializer的工作就是遍历这个复杂的Protobuf结构将其扁平化地组装成我们之前看到的那种JSON对象。我们的自定义序列化器需要介入这个过程。在Canal的Kafka生产者配置中我们可以指定key.serializer和value.serializer。通常我们只关心value.serializer因为消息内容在这里。我们需要实现org.apache.kafka.common.serialization.Serializer接口并在Canal的instance.properties配置文件中进行配置。注意Canal的Kafka输出模块在canal.properties中有一个canal.mq.serializer的配置项但其默认实现和直接配置Kafka Producer的Serializer略有不同。经过测试为了获得最大的灵活性和对Kafka原生特性的支持我们选择直接配置Kafka Producer的方式即设置kafka.key.serializer和kafka.value.serializer。这要求我们对Canal的Kafka生产者初始化代码有一定的了解或者使用更高版本的Canal1.1.5其对这种配置方式支持得更好。4. 自定义JSON序列化器的设计与实现我们设计了一个名为CustomCanalJsonSerializer的类。核心思路是实现Kafka的SerializerMessage接口Canal封装了一层Message对象里面包含了CanalEntry.Entry的列表在serialize方法中完成格式转换。4.1 定义目标JSON格式首先我们和下游团队共同敲定了两种目标格式模板格式A精简扁平式适用于大多数业务服务消费。{ op: c, // 操作类型: c创建, u更新, d删除 db: test_db, tb: user, ts: 1696147200123, // 事件发生时间戳毫秒 data: { // 变更后的数据对于DELETE此对象为空或包含删除前的数据 id: 1, // 注意根据sqlType尝试转换类型 name: test, createTime: 1696147200000 // 时间字段转换为时间戳 }, old: { // 仅UPDATE操作存在记录变更前的值 name: old_name } }格式B原生增强式适用于需要最大信息量的数据平台或审计场景。它在默认格式基础上做了优化比如时间戳标准化、字段名风格转换可选。4.2 核心转换逻辑实现serialize方法是核心。其步骤如下提取CanalEntry.Entry从输入的CanalMessage对象中获取Entry列表。通常一个Message包含多个Entry我们需要遍历处理。过滤与判断只处理EntryType为ROWDATA的Entry。忽略事务开始/结束等其他类型的Entry。解析RowChange将Entry的StoreValue反序列化成RowChange对象。构建基础信息从Entry的Header中提取databaseName,tableName,executeTime事件执行时间。处理行数据遍历RowChange的rowDatas。根据eventType决定操作类型(op)。INSERT取afterColumns构建data对象。UPDATE取afterColumns构建data对象同时取beforeColumns中updatedtrue的列构建old对象。DELETE取beforeColumns构建data对象记录被删除的数据或者将data置空仅通过op标识删除。字段值转换这是格式定制的关键。命名风格转换实现一个FieldNamingStrategy比如CAMEL_CASE_TO_UNDERSCORE或UNDERSCORE_TO_CAMEL_CASE。可以通过配置开关控制。类型转换根据Column的sqlType尝试将字符串类型的value转换为更合适的JSON类型。例如sqlType为BIGINT或INTEGER则将value转为Long或IntegersqlType为TIMESTAMP或DATETIME则将字符串“2023-10-01 12:00:00”转换为毫秒时间戳Long。对于无法确定或转换失败的情况保持原字符串。字段过滤提供一个可配置的白名单或黑名单只输出需要的字段或者排除掉如is_deleted这类敏感或不必要的字段。组装JSON并序列化使用如Jackson或Gson库将构建好的Java对象Map或自定义POJO序列化成JSON字符串再转换为UTF-8字节数组返回。4.3 关键配置与代码片段在instance.properties中配置Kafka Producer使用我们的序列化器# Kafka基本配置 kafka.bootstrap.servers localhost:9092 kafka.acks all kafka.retries 3 # 指定自定义序列化器 kafka.key.serializer org.apache.kafka.common.serialization.StringSerializer kafka.value.serializer com.yourcompany.canal.serializer.CustomCanalJsonSerializer # 自定义序列化器的参数通过ProducerConfig传递 custom.serializer.format.type FORMAT_A custom.serializer.field.naming.strategy CAMEL_CASE custom.serializer.convert.timestamp true custom.serializer.include.fields id,name,createTime序列化器实现中需要在configure方法中读取这些自定义参数public class CustomCanalJsonSerializer implements SerializerMessage { private String formatType; private ObjectMapper objectMapper; Override public void configure(MapString, ? configs, boolean isKey) { this.formatType (String) configs.get(custom.serializer.format.type); this.objectMapper new ObjectMapper(); // 根据formatType和配置设置objectMapper的属性如命名策略、日期格式 if (CAMEL_CASE.equals(configs.get(custom.serializer.field.naming.strategy))) { objectMapper.setPropertyNamingStrategy(PropertyNamingStrategies.LOWER_CAMEL_CASE); } // 配置时间戳转换 if (true.equals(configs.get(custom.serializer.convert.timestamp))) { objectMapper.configure(SerializationFeature.WRITE_DATES_AS_TIMESTAMPS, true); } } Override public byte[] serialize(String topic, Message data) { if (data null) return null; ListCanalEntry.Entry entries data.getEntries(); ListMapString, Object formattedMessages new ArrayList(); // ... 遍历entries进行上述转换逻辑填充formattedMessages ... try { return objectMapper.writeValueAsBytes(formattedMessages); } catch (JsonProcessingException e) { throw new SerializationException(Error serializing Canal message to JSON, e); } } // ... close方法 ... }实操心得在类型转换时要特别注意NULL值的处理。Canal中isNull为true时value可能是空字符串。在转换时应直接输出JSON的null而不是字符串“”或“null”。另外对于数值类型转换前要做好异常捕获防止因脏数据导致整个序列化失败。5. 部署、测试与上下游联调开发完成后部署是关键一步。我们采用Docker部署Canal Server将打包好的自定义序列化器JAR包及其依赖挂载到Canal容器的lib目录下。然后修改对应instance的配置文件。部署步骤打包将CustomCanalJsonSerializer及其依赖如Jackson打包成一个独立的JAR比如canal-custom-serializer-1.0.0.jar。放置在Canal Server的部署目录例如/opt/canal-server下创建ext-lib目录Canal会自动加载此目录下的JAR将我们的JAR包放进去。或者直接放入lib目录但要注意版本冲突。配置修改目标实例的instance.properties如上节所示配置序列化器和相关参数。重启重启Canal Server实例。观察日志确保没有ClassNotFoundException等错误并且Kafka Producer初始化时加载了我们的序列化器。测试验证测试分两步走单元测试编写单元测试模拟CanalEntry.Entry数据调用序列化器的serialize方法验证输出的JSON字符串是否符合预期格式特别是边界情况如NULL值、特殊字符、所有操作类型。集成测试启动一个测试用的MySQL和Kafka。配置Canal连接测试库并指向测试Kafka集群。在测试库的表上执行INSERT、UPDATE、DELETE操作。使用Kafka控制台消费者或编写一个简单消费者从对应的Topic拉取消息肉眼观察格式是否正确。关键点验证时间戳转换、字段名风格转换、类型转换字符串数字转JSON数字是否生效。上下游联调这是最体现价值的一环。我们邀请下游几个核心消费方的开发者一起联调。提供“数据契约”文档将我们定义好的两种JSON格式格式A和格式B以文档形式明确下来包括每个字段的含义、类型、可能的值。这相当于一份API合同。共享测试Topic让下游团队从这个Topic消费数据用他们自己的反序列化逻辑如Java的JacksonJsonPropertyGo的json.Unmarshal来解析。收集反馈并迭代下游可能会提出新的需求比如“希望把变更的字段列表单独提出来”、“希望把主键值放在更顶层”。我们在评估后通过增加序列化器的配置项来满足合理的需求。这个过程可能来回几次直到双方都觉得格式“顺手”为止。注意事项联调阶段一定要做好向后兼容。在格式版本升级时可以考虑使用不同的Kafka Topic或者在新格式消息中增加一个version字段让下游消费者有能力同时处理新旧格式并逐步迁移。切忌直接修改现有Topic的消息格式这会导致下游所有正在运行的消费者立即崩溃。6. 性能考量、监控与优化建议改了格式性能怎么样这是必须回答的问题。我们从序列化开销、网络开销、消费端解析开销三个方面做了评估。序列化开销自定义序列化器比简单的StringSerializer直接输出预格式字符串肯定要多一些CPU计算主要是类型判断和转换。我们使用JMH做了基准测试在单线程下处理一条典型消息约20个字段自定义序列化耗时比默认方式多出约15%-30%。这个开销在可接受范围内因为Canal本身是I/O密集型读binlog、网络发送CPU不是主要瓶颈。如果担心可以启用Jackson的Smile二进制JSON格式或者对序列化器对象做线程本地缓存(ThreadLocalObjectMapper)来减少对象创建开销。网络与存储开销这是收益最明显的地方。通过字段过滤和精简我们的一条消息平均大小从原来的约2KB减少到约800字节格式A。这意味着Kafka的吞吐量间接提升了下游消费的网络传输和解析速度也会更快。对于按流量计费的云服务这直接节省了成本。消费端解析开销下游反馈由于字段名符合编程习惯、类型直接可用不用再String转Long他们的反序列化代码更简洁解析性能提升了约40%。监控方案格式改造上线后监控至关重要。Canal Server监控关注canal.instance.network.error、canal.instance.download.error等错误计数确保解析和发送本身正常。Kafka Producer监控通过JMX或Kafka自身指标监控record-send-rate、record-error-rate、request-latency-avg。自定义序列化器的异常会被计入record-error-rate。消息格式监控我们编写了一个简单的监控消费者订阅所有相关的Topic对每条消息进行格式校验如JSON语法、必填字段是否存在、类型是否正确将格式错误的消息和计数打到日志和监控系统如Prometheus中。业务端监控与下游约定在他们的消费逻辑里对解析失败的消息进行计数和告警并最好能将错误消息原文落盘方便排查。优化建议批量序列化Canal的Message可能包含多个Entry。我们的序列化器是逐个Entry处理再组装成数组。可以考虑更高效的流式JSON生成方式直接构建一个大的JSON数组减少中间对象的创建。条件性字段包含通过配置实现更动态的字段输出。例如只有当字段值非空时才输出或者根据数据库表名动态选择不同的字段映射规则。Schema Registry集成进阶对于非常复杂的格式管理和演进可以考虑将Avro、Protobuf与Schema Registry结合。Canal输出Avro格式到KafkaSchema Registry保证兼容性。但这套方案更重适用于大型数据平台。7. 常见问题排查与实战经验录在实际操作中我们遇到了不少问题这里总结几个典型的问题一序列化器JAR包加载失败Canal启动报ClassNotFoundException。排查首先检查JAR包是否放到了Canal Server的classpath中。ext-lib目录是最推荐的位置。其次检查JAR包是否有依赖缺失。可以用jar tf your-serializer.jar查看包内内容或用mvn dependency:copy-dependencies打出所有依赖包一并放入ext-lib。解决确保所有依赖都被正确加载。一个稳妥的办法是使用maven-shade-plugin或spring-boot-maven-plugin打一个包含所有依赖的fat jar。问题二消息发送到Kafka成功但下游解析时发现字段值为null而数据库中明明有值。排查检查Canal的instance.properties中是否配置了canal.instance.filter.query.dml false默认是false如果为true会过滤掉DML语句的binlog。更重要的是检查自定义序列化器中的字段名转换逻辑。很可能是因为数据库字段名是user_name转换成了userName但下游反序列化时期待的字段名是user_name导致不匹配。解决核对序列化器配置的field.naming.strategy与下游反序列化时使用的命名策略。最好双方使用相同的JSON库如Jackson并配置相同的PropertyNamingStrategy。提供一个配置开关让字段名风格可配。问题三UPDATE操作时old字段包含了所有字段而不仅仅是变更的字段。排查这是对Canal数据模型理解有误。Canal的RowData中beforeColumns列表里只有那些updatedtrue的列才表示该列在本次UPDATE中发生了变化。我们的序列化器在构建old对象时应该只选择updatedtrue的列。解决修改序列化逻辑在遍历beforeColumns时判断column.getUpdated()是否为true只有为true的才放入old对象中。问题四时间字段转换后下游收到的时间戳是负数或明显不对。排查时区问题。Canal从binlog中解析出的时间是MySQL服务器的时间字符串。我们的序列化器在将其转换为时间戳时必须明确时区。如果MySQL用的是UTC8而转换代码默认用了UTC就会差8小时。解决在序列化器中使用数据库的时区可以从Canal配置或连接参数中获取来解析日期字符串。例如LocalDateTime.parse(str, formatter).atZone(ZoneId.of(Asia/Shanghai)).toInstant().toEpochMilli()。问题五Kafka消息量巨大序列化器成为瓶颈CPU使用率高。排查使用Profiling工具如Arthas的profiler命令对Canal Server进行采样查看热点是否在自定义序列化器的serialize方法中。解决检查是否在每次serialize调用中都创建了新的ObjectMapper实例将其改为静态成员或ThreadLocal缓存。检查类型转换逻辑是否过于复杂尝试简化或对常见的sqlType做快速路径优化。考虑是否过滤了足够多的无用字段减少操作的数据量是最根本的优化。如果单实例压力仍然很大可以考虑对Canal实例进行水平拆分一个实例只负责同步一部分表分散压力。这次与Canal和Kafka打交道的经历让我深刻体会到数据管道中的每一个环节其输出格式都是一种API契约。设计一个考虑周全、灵活可配的契约不仅能提升下游的开发体验更能为整个数据流的高效、稳定运行打下坚实的基础。格式定制看似是个“细枝末节”的改造实则牵一发而动全身需要我们对数据源、工具链和业务需求都有深入的理解。