Kafka 2.6配置问询:如何让消费者仅消费服务端正确序列化的合法记录
Kafka 2.6 + Confluent Avro 组件坏记录拦截配置方案
核心逻辑
要实现仅消费服务端合法序列化的记录,需要结合服务端Broker写入校验和消费者反序列化容错两层配置,从源头拦截非法Avro记录,避免坏记录流入消费逻辑。
第一步:服务端Broker配置(必做,从源头拦截非法消息)
Kafka 2.6搭配Confluent Avro序列化组件时,可以在Broker端配置Schema校验规则,只有符合注册Schema的Avro消息才会被持久化到Topic,非法消息直接在生产写入阶段就被拒绝。
需要在server.properties中添加以下配置:
- 配置Schema Registry访问地址:
confluent.schema.registry.url=http://你的SchemaRegistry实例地址:8081 - 对目标Topic启用Schema校验,将
[Topic名称]替换为实际的业务Topic名:confluent.topic.[Topic名称].schema.validation.enable=true - 配置全量校验策略(同时校验消息Key和Value的Schema合法性):
confluent.topic.[Topic名称].schema.validation.strategy=ALL
注意:修改Broker配置后需要滚动重启Broker集群生效,配置完成后非法Avro消息在生产者发送时就会返回写入失败,不会进入Topic队列。
第二步:消费者端容错配置(补充,避免漏网的坏消息中断消费)
如果暂时无法修改服务端配置,可以在消费者侧配置反序列化容错机制,遇到坏记录时自动跳过,不会抛出异常中断消费流程。
消费者consumer.properties需要修改以下配置:
- 使用Confluent官方Avro反序列化器:
- Key反序列化配置:
key.deserializer=io.confluent.kafka.serializers.KafkaAvroDeserializer - Value反序列化配置:
value.deserializer=io.confluent.kafka.serializers.KafkaAvroDeserializer
- Key反序列化配置:
- 配置Schema Registry访问地址:
schema.registry.url=http://你的SchemaRegistry实例地址:8081 - 启用坏记录自动跳过逻辑:
- 如果使用SpecificRecord类型消费,开启:
specific.avro.reader=true - 反序列化失败不抛出异常:
fail.on.invalid.message=false - 配置反序列化失败策略为自动跳过:
deserialization.failure.strategy=org.apache.kafka.common.serialization.IgnoreDeserializationFailureStrategy
- 如果使用SpecificRecord类型消费,开启:
- 可选配置:开启坏记录日志打印,方便排查问题:
log.deserialization.failures=true,反序列化失败时会打印坏消息的偏移量、所属Topic等信息到服务日志。
可选增强:自定义拦截器实现灵活处理
如果需要更复杂的坏消息处理逻辑(比如把坏消息投递到死信队列留痕),可以自定义消费拦截器实现ConsumerInterceptor接口,在consume方法中判断消息是否为反序列化失败的null值,将合法消息返回给业务逻辑,非法消息统一转存到死信Topic即可。
内容的提问来源于stack exchange,提问作者sowsls
相关产品推荐
相关产品推荐

