自定义Kafka反序列化器忽略空消息,如何让空消息进入deserialize方法?
Kotlin Spring Boot自定义Kafka Deserializer处理空消息问题
问题描述
在使用Kotlin Spring Boot实现自定义Kafka Deserializer时,带有null值的消息会被完全忽略,无法进入deserialize方法,需要让空消息能进入该方法处理。
现有配置
application.yaml配置:
spring: kafka: consumer: value-deserializer: mypackage.CustomDeserializer key-deserializer: org.apache.kafka.common.serialization.StringDeserializer
自定义Deserializer代码
class CustomDeserializer : Deserializer<Result<MyData?>> { private val log = logger() override fun deserialize(topic: String?, data: ByteArray?): Result<MyData?> { // 反序列化逻辑 } }
解决办法
默认Kafka Consumer会过滤value为null的消息,不会传递给自定义反序列化器。只需在消费者配置中添加allow-null-values: true参数,即可让空消息进入deserialize方法:
spring: kafka: consumer: allow-null-values: true value-deserializer: mypackage.CustomDeserializer key-deserializer: org.apache.kafka.common.serialization.StringDeserializer
添加该配置后,当消息value为null时,deserialize方法的data参数会是null,你可以在方法内针对性处理,比如返回Result.success(null)或其他符合业务需求的结果。
注意:若使用的Kafka版本低于2.0,
allow-null-values参数不存在,此时可通过拦截ConsumerRecord的方式在反序列化前捕获空消息处理,但这种方式不如配置参数简洁。
内容的提问来源于stack exchange,提问作者T. Aris
相关产品推荐
相关产品推荐

