Spring Cloud Stream批处理部分失败时处理成功记录与失败消息方案
使用函数式风格的Spring Cloud Stream开发批处理流应用时,完全可以实现「批次内成功记录正常流转下游、失败记录单独重试/投递DLQ」的需求,你当前遇到的「失败位置前的成功消息不发送」问题,是框架默认批处理语义和代码实现方式共同导致的。
Spring Cloud Stream Kafka绑定器默认对批处理采用全量成功语义:只要业务函数(即你的apply方法)执行过程中抛出异常未正常返回,函数执行过程中构造的所有待下发输出消息都会被丢弃,整个批次会被判定为处理失败,交由你配置的DefaultErrorHandler处理。
你的现有实现存在两个直接触发问题的逻辑:
- 遍历批次记录时,只要某条记录处理失败就直接抛出
BatchListenerFailedException,中断整个函数执行,已经加入output列表的成功记录没有机会被框架发送到下游 - 未开启单条记录确认配置,容器默认只有整个批次处理成功才会提交偏移、下发输出消息,不支持部分成功的处理场景
方案1:业务层拆分成功/失败记录,不中断函数执行
这是改造成本最低、最适配现有代码结构的方案:
- 遍历批次时,单条记录的处理异常单独捕获,不向上抛出中断整个函数流程
- 处理成功的记录正常加入输出列表,保证函数正常返回时这部分消息可以被正常投递到下游
- 单独收集处理失败的记录,在所有记录遍历完成后,复用现有重试/DLQ逻辑处理这部分失败记录
参考改造后的业务代码:
@Override public List<Message<Context>> apply(Message<List<Context>> listMessage) { List<Message<Context>> output = new ArrayList<>(); List<RecordFailureCtx> failureCtxList = new ArrayList<>(); IntStream.range(0, listMessage.getPayload().size()).forEach(index -> { try { Record<Context> record = Record.fromBatch(listMessage, index); output.add(MessageBuilder.withPayload(record.getValue()).build()); // 模拟单条记录处理失败场景 if (index == listMessage.getPayload().size() - 1) { throw new TransientError("offset " + record.getOffset() + " failed", new RuntimeException()); } } catch (Exception e) { // 收集失败上下文,不抛出异常中断批次 failureCtxList.add(new RecordFailureCtx(listMessage, index, e)); } }); // 批量处理收集到的失败记录:执行配置好的重试策略,重试耗尽后投递DLQ if (!failureCtxList.isEmpty()) { failureHandler.handleFailures(failureCtxList); } // 正常返回成功记录,触发下游下发 return output; }
*如果要完全复用你现有自定义的退避策略、异常分类、DLQ头构造逻辑,不需要重复实现重试代码,可以把失败记录发送到内部专用的错误topic,绑定到你已经配置好的错误处理链路即可。
方案2:调整绑定器配置,开启容器层面的批次部分成功支持
Spring Kafka 2.8.0版本的DefaultErrorHandler本身支持批次处理时仅定位失败记录重试/投递DLQ,不影响已处理成功的记录,只需要补充对应配置:
- 开启单条记录确认能力,允许容器提交已处理成功记录的偏移、下发对应输出
- 关闭批处理事务,避免异常触发整个批次的输出回滚
需要在YAML中补充的配置:
spring: cloud: stream: kafka: bindings: batchFunc-in-0: consumer: # 开启单条记录粒度的确认 ackEachRecord: true # 关闭批次级事务,避免异常回滚全部输出 transactional: false
*使用该方案需要注意:抛出BatchListenerFailedException的时机必须是单条记录处理失败的当下,不要等所有记录遍历完成再抛出异常,否则框架无法准确识别已处理成功的记录边界。
方案3:单条记录粒度接入重试模板
你当前已经配置spring.cloud.stream.default.consumer.maxAttempts=1关闭了框架内置的批处理重试,可以直接在单条记录的catch块中注入自定义配置的RetryTemplate,按照你定义的退避策略、异常分类规则重试单条失败记录,重试耗尽后直接调用已构造的DeadLetterPublishingRecoverer投递DLQ。该方案灵活度最高,可以针对不同业务场景定制单条记录的处理逻辑,完全不会影响批次内其他记录的正常流转。
- 你自定义的
CommonErrorHandler是监听容器层面的处理器,只有业务函数抛出异常时才会触发,此时函数已经中断,输出列表的消息已经被丢弃,因此无法实现成功记录的先下发 - 你当前代码中是在遍历到最后一条记录触发异常时才抛出
BatchListenerFailedException,此时虽然前面的记录已经加入输出列表,但函数没有正常返回,框架不会触发下游发送逻辑
内容的提问来源于stack exchange,提问作者Yosi Bronsberg

