Spring Kafka 3.1.3中如何通过ErrorHandlingDeserializer实现自定义日志与指标?
解决Spring Kafka 3.1.3中毒丸消息的自定义日志与指标统计需求
针对你的场景,不要重写handleRemaining()方法,这个方法是用来处理轮询批次中剩余未处理的记录,无法精准定位单条失败的毒丸消息。正确的做法是通过DefaultErrorHandler搭配自定义的ConsumerRecordRecoverer,针对单条失败记录做精准处理。
实现步骤
1. 自定义毒丸消息恢复处理器
实现ConsumerRecordRecoverer接口,在其中完成毒丸消息的日志记录和Micrometer指标统计:
import org.slf4j.Logger; import org.slf4j.LoggerFactory; import org.springframework.kafka.listener.ConsumerRecordRecoverer; import org.springframework.kafka.support.serializer.DeserializationException; import io.micrometer.core.annotation.Counted; import io.micrometer.core.instrument.MeterRegistry; import org.apache.kafka.clients.consumer.ConsumerRecord; public class PoisonPillRecoveryHandler implements ConsumerRecordRecoverer { private static final Logger log = LoggerFactory.getLogger(PoisonPillRecoveryHandler.class); private final MeterRegistry meterRegistry; public PoisonPillRecoveryHandler(MeterRegistry meterRegistry) { this.meterRegistry = meterRegistry; } @Override public void accept(ConsumerRecord<?, ?> record, Exception exception) { if (exception instanceof DeserializationException deserializationEx) { // 提取毒丸消息元数据 String topic = record.topic(); int partition = record.partition(); long offset = record.offset(); String key = String.valueOf(record.key()); // 自定义日志格式 log.error("[毒丸消息拦截] 主题: {}, 分区: {}, 偏移量: {}, 键值: {}, 异常原因: {}", topic, partition, offset, key, deserializationEx.getCause().getMessage()); // 生成Micrometer自定义指标 meterRegistry.counter("kafka.poison_pill.total", "topic", topic, "partition", String.valueOf(partition)) .increment(); } } }
2. 配置DefaultErrorHandler
通过Bean配置将自定义恢复处理器关联到DefaultErrorHandler,并指定仅对反序列化异常直接跳过(不重试):
import org.springframework.context.annotation.Bean; import org.springframework.context.annotation.Configuration; import org.springframework.kafka.listener.DefaultErrorHandler; import org.springframework.kafka.support.serializer.DeserializationException; import io.micrometer.core.instrument.MeterRegistry; @Configuration public class KafkaErrorConfig { @Bean public DefaultErrorHandler kafkaErrorHandler(MeterRegistry meterRegistry) { PoisonPillRecoveryHandler recoveryHandler = new PoisonPillRecoveryHandler(meterRegistry); DefaultErrorHandler errorHandler = new DefaultErrorHandler(recoveryHandler); // 标记反序列化异常为无需重试,直接交给恢复处理器处理 errorHandler.addNotRetryableException(DeserializationException.class); return errorHandler; } }
原理说明
- 你的
ErrorHandlingDeserializer配置已经正确,当遇到反序列化失败时,会抛出DeserializationException,并将原始的ConsumerRecord携带在异常上下文里。 DefaultErrorHandler会捕获该异常,由于我们标记了它为无需重试,会直接将失败的ConsumerRecord和异常传递给自定义的PoisonPillRecoveryHandler。- 在恢复处理器中,你可以精准获取单条毒丸消息的元数据(主题、分区、偏移量等),完成自定义日志和指标统计,同时跳过该消息不影响后续消息处理。
内容的提问来源于stack exchange,提问作者Rod
相关产品推荐
相关产品推荐

