Kafka Streams使用Header存schema标识的Serdes调用及类转换异常问题
问题解答
1. 关于优先调用带headers的序列化方法的说明
Kafka Streams从2.7.0版本开始,默认已经会优先调用带Headers参数的serialize(String topic, Headers headers, T data)/deserialize方法,仅当序列化器没有实现该重载方法时,才会降级调用无headers的版本。
你遇到的内部调用无headers方法的场景,是KTable外键关联的框架层面限制:部分版本的KTable外键join内部处理器存在未透传Headers的问题,无法通过普通配置直接修改。
2. transformValues方案报错的修复方法
你遇到的ClassCastException核心原因是transformValues操作没有显式指定输出结果的Materialized配置,Kafka Streams默认继承了上游KTable的Serdes.ByteArray()作为状态存储的序列化器,但是你transform之后输出的value是MySpecificClass类型,用ByteArraySerializer序列化该类型自然会抛出类型转换异常。
修复只需给transformValues方法传入匹配输出类型的Materialized配置即可:
final KTable<Long, MySpecificClass> myTable = builder.table( "my-topic", Consumed.with(Serdes.Long(), Serdes.ByteArray()) ) .transformValues(MyDeserializerTransformer::new, // 这里传入匹配输出类型的Serde,仅需要实现基础序列化逻辑,不需要处理Headers Materialized.with(Serdes.Long(), mySpecificClassSerde));
注意这里的自定义Serde是用于内部状态存储的本地序列化,不需要发送到Kafka broker,也不需要写入带schema标识的Headers,实现成本很低。
3. 其他可行的规避方案
- 自定义Serde兼容无headers方法:如果你的自定义Avro序列化器仅在写入Kafka的时候需要Headers写入schema标识,可以给无headers的
serialize重载方法加兼容逻辑:内部序列化场景下直接返回不带schema标识的序列化结果,不需要抛出异常即可适配框架的内部调用逻辑。 - 预转换数据到新topic:如果以上方案都不符合你的业务要求,可以先运行一个轻量的Streams任务,将原topic的Avro数据读取后,重新序列化为schema标识写入消息体的兼容格式写入新topic,后续的外键join任务直接读取新topic的数据即可,该方案兼容性最好,适合长期稳定运行的生产场景。
4. 关于二级状态存储的说明
不需要使用二级状态存储,上述方案已经可以解决所有问题,额外引入二级状态存储只会增加逻辑复杂度和运行开销。
内容的提问来源于stack exchange,提问作者MMEK
相关产品推荐
相关产品推荐

