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(Kafka偏移量标识的是下一条要消费的记录位置)
- 使用
commitSync同步提交偏移量,确保偏移量提交成功后再继续消费下一批次,避免重复消费或漏消费
内容的提问来源于stack exchange,提问作者Shahul Basha
相关产品推荐
相关产品推荐

