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

Spring Kafka批量监听器自定义错误处理器实现咨询

Spring Kafka 自定义批量错误处理器实现

需求说明

  • 批量监听器处理批次时,若某条记录触发异常,通过开关控制两种处理逻辑:
    • 开关shouldRetry为true:无限重试当前批次,直至人工介入
    • 开关shouldRetry为false:跳过当前批次,直接处理下一批次

代码实现(修正后)

import org.apache.kafka.clients.consumer.Consumer;
import org.apache.kafka.clients.consumer.ConsumerRecords;
import org.apache.kafka.clients.consumer.OffsetAndMetadata;
import org.apache.kafka.common.TopicPartition;
import org.slf4j.Logger;
import org.slf4j.LoggerFactory;
import org.springframework.kafka.listener.DefaultErrorHandler;
import org.springframework.kafka.listener.MessageListenerContainer;

import java.util.HashMap;
import java.util.Map;

public class CustomBatchErrorHandler extends DefaultErrorHandler {

    private static final Logger LOG = LoggerFactory.getLogger(CustomBatchErrorHandler.class);
    private final boolean shouldRetry;

    public CustomBatchErrorHandler() {
        super();
        this.shouldRetry = getFlag(); // 替换为实际开关获取逻辑
    }

    @Override
    public void handleBatch(Exception thrownException, ConsumerRecords<?, ?> data, Consumer<?, ?> consumer,
                            MessageListenerContainer container, Runnable invokeListener) {
        if (shouldRetry) {
            LOG.warn("因异常重试当前批次", thrownException);
            super.handleBatch(thrownException, data, consumer, container, invokeListener);
        } else {
            LOG.warn("跳过当前批次,继续处理下一批次", thrownException);
            // 提交当前批次所有分区的最后偏移量,确认已处理(跳过)
            Map<TopicPartition, OffsetAndMetadata> offsetMap = new HashMap<>();
            data.partitions().forEach(partition -> {
                long lastRecordOffset = data.records(partition).get(data.records(partition).size() - 1).offset();
                offsetMap.put(partition, new OffsetAndMetadata(lastRecordOffset + 1));
            });
            consumer.commitSync(offsetMap);
        }
    }

    // 示例:获取开关状态的方法,根据实际业务场景实现
    private boolean getFlag() {
        // 可替换为从配置文件、配置中心或数据库读取开关状态
        return false;
    }
}

关键逻辑说明

  • 当开关为false时,核心操作是提交当前批次所有分区的最后偏移量:
    1. 遍历当前批次包含的所有分区,获取每个分区最后一条记录的偏移量
    2. 提交时偏移量需+1(Kafka偏移量标识的是下一条要消费的记录位置)
    3. 使用commitSync同步提交偏移量,确保偏移量提交成功后再继续消费下一批次,避免重复消费或漏消费

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

相关产品推荐
方舟 Agent Plan

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

最近更新时间:2026.06.24 07:34:57