You need to enable JavaScript to run this app.
优惠活动
大模型
产品
解决方案
定价
更多

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;
    }
}

原理说明

  1. 你的ErrorHandlingDeserializer配置已经正确,当遇到反序列化失败时,会抛出DeserializationException,并将原始的ConsumerRecord携带在异常上下文里。
  2. DefaultErrorHandler会捕获该异常,由于我们标记了它为无需重试,会直接将失败的ConsumerRecord和异常传递给自定义的PoisonPillRecoveryHandler。
  3. 在恢复处理器中,你可以精准获取单条毒丸消息的元数据(主题、分区、偏移量等),完成自定义日志和指标统计,同时跳过该消息不影响后续消息处理。

内容的提问来源于stack exchange,提问作者Rod

相关产品推荐
方舟 Agent Plan

超全模态模型 × Harness 升级,最新支持 Deepseek-V4.1-Flash、GLM-5.3 系列、Doubao-Seedream-5.0-pro、Kimi-K3 (部分), 限时 9.9 元起

最近更新时间:2026.06.24 06:11:05