Spring Cloud Stream Kafka Binder批量消费重试3次及DLQ投递实现问题
问题根源说明
- 重试次数超过设置值的原因:
SeekToCurrentBatchErrorHandler默认未配置重试耗尽后的恢复逻辑,重试次数用完后会继续拉取同一批消息重复处理,看起来就是次数无限超过设置值。 - 原有监听器中
deliveryAttempt判断失效的原因:该头是单条消费场景下的重试计数头,批量消费模式下框架不会为整批消息维护计数,你拿到的永远是默认值1,自然走不到发DLQ的逻辑。 - 无法导入
RetryingBatchErrorHandler的原因:该类是Spring Kafka 2.3.7版本才新增的,你当前Spring Cloud版本依赖的Spring Kafka版本低于该要求,所以找不到类。
适配当前版本的解决方案(无需升级依赖)
直接基于你正在用的SeekToCurrentBatchErrorHandler配置批量恢复逻辑即可,同时满足整批最多重试3次、失败自动发DLQ的需求:
1. 修改配置类代码
@Configuration public class ConsumerConfig { // 注入你自己实现的DLQ发送工具类 @Autowired private DlqPublisher publisher; @Bean public ListenerContainerCustomizer<AbstractMessageListenerContainer<?, ?>> customizer() { return (container, dest, group) -> { container.getContainerProperties().setAckOnError(false); SeekToCurrentBatchErrorHandler errorHandler = new SeekToCurrentBatchErrorHandler(); // FixedBackOff第二个参数是重试次数,设置为2代表初始调用1次+重试2次,总共执行3次 errorHandler.setBackOff(new FixedBackOff(0L, 2L)); // 配置重试耗尽后的恢复逻辑:整批发送到DLQ errorHandler.setBatchRecoverer((records, exception) -> { log.error("批量消费重试3次全部失败,消息条数:{},异常信息:", records.count(), exception); // 遍历批量消息组装成你需要的格式发送到DLQ List<Event> dlqEvents = new ArrayList<>(); records.forEach(record -> dlqEvents.add((Event) record.value())); publisher.publishToDlq(EventType.UPDATE_FAILED, dlqEvents, exception.getMessage()); // 如果你用的是手动ACK模式,这里需要手动提交偏移量,避免重复拉取 // records.consumer().commitSync(); }); container.setBatchErrorHandler(errorHandler); }; } }
2. 简化监听器代码
不需要再手动捕获异常判断重试次数,处理失败直接抛出异常即可,框架会自动处理重试和DLQ投递:
@StreamListener(ActivityChannel.INPUT_CHANNEL) public void handleActivity(List<Message<Event>> messages, @Header(name = KafkaHeaders.ACKNOWLEDGMENT) Acknowledgment acknowledgment) { log.info("Received activity message with message length {}", messages.size()); nodeConfigActivityBatchProcessor.processNodeConfigActivity(messages); acknowledgment.acknowledge(); log.debug("Processed activity message {} successfully!!", messages.size()); }
升级依赖使用RetryingBatchErrorHandler的方案(可选)
如果要使用RetryingBatchErrorHandler,先将Spring Kafka版本升级到2.3.7及以上,或升级Spring Cloud版本到Hoxton.SR8及以上,配置逻辑和上述方案一致,仅需要把SeekToCurrentBatchErrorHandler替换为RetryingBatchErrorHandler即可。
内容的提问来源于stack exchange,提问作者APK
相关产品推荐
相关产品推荐

