Kafka SourceRecord使用自定义对象作为sourceOffset查询返回null如何解决
你的猜测是正确的。
Kafka Connect的偏移量存储默认使用JsonConverter做序列化和反序列化,仅原生支持字符串、数值、布尔值、null,以及由这些类型组成的Map、List结构。你传入自定义对象时,默认序列化逻辑无法正确解析对象结构,会导致偏移量存储异常,查询时自然返回null。
解决方案
方案1(最推荐,无额外配置)
不要直接传入自定义对象作为sourceOffset的value,先将自定义对象转换为符合原生支持类型的Map<String, Object>结构,所有嵌套值都使用基础类型即可。
示例代码:
// 假设你的自定义偏移量对象为CustomOffset CustomOffset customOffset = getMyOffset(); // 转换为原生类型组成的Map Map<String, Object> offsetMap = Map.of( "userId", customOffset.getUserId(), "lastSyncTime", customOffset.getLastSyncTime(), "pageNum", customOffset.getPageNum() ); // 传入SourceRecord SourceRecord record = new SourceRecord(partition, offsetMap, topic, null, valueSchema, value);
查询到偏移量Map后,再手动转为自定义对象即可,无需修改任何Kafka Connect配置,兼容性最好。
方案2(需要修改Worker配置,适合特殊场景)
如果必须直接存储自定义对象,你需要自定义序列化/反序列化逻辑,修改Kafka Connect Worker的全局配置:
- 自定义支持你的对象类型的JSON序列化器/反序列化器(可基于Jackson实现,添加对应的类型映射或多态序列化注解)
- 修改Worker配置文件中的
offset.storage.converter参数,指定为你自定义的Converter实现类,同时配置对应的序列化/反序列化参数
该方案会影响整个Worker集群的偏移量存储逻辑,不推荐普通场景使用。
内容的提问来源于stack exchange,提问作者alrevuelta
相关产品推荐
相关产品推荐

