如何使用Mono实现带状态的失败条目最多N次重试发送逻辑
最优实现方案(基于Reactor expand 操作符)
你当前的场景非常适合用Reactor提供的expand操作符实现,它是专门为迭代生成下一个Publisher的场景设计的,不需要额外维护可变状态,也不需要靠抛出异常触发重试,完全符合响应式编程规范,和你原有阻塞逻辑的语义完全对齐。
第一步:适配阻塞方法为响应式Mono
如果你的sendItems本身是阻塞调用,先适配为非阻塞Mono,注意要指定调度器隔离阻塞操作:
// 若sendItems本身已经是异步返回Mono,可跳过这步直接使用原有方法 private Mono<List<String>> sendsItemsReactive(List<String> items) { return Mono.fromCallable(() -> sendsItems(items)) // 阻塞操作放到boundedElastic调度器,避免占用业务EventLoop线程 .subscribeOn(Schedulers.boundedElastic()); }
第二步:核心重试逻辑实现
/** * 发送条目并仅重试失败项 * @param initialItems 初始待发送条目 * @param maxAttempts 最大尝试次数(和你原有阻塞逻辑对齐,比如传3表示最多尝试发送3次) * @return 最终发送完成信号 */ public Mono<Void> sendWithRetries(List<String> initialItems, int maxAttempts) { return sendsItemsReactive(initialItems) // expand接收上一次返回的失败条目列表,生成下一次要执行的发送操作 .expand(failedItems -> { // 没有失败条目直接返回空Mono,终止迭代 if (failedItems.isEmpty()) { return Mono.empty(); } // 可以在这里添加重试延迟,比如重试前等待1秒: // return Mono.delay(Duration.ofSeconds(1)).then(sendsItemsReactive(failedItems)); return sendsItemsReactive(failedItems); }) // 限制最大尝试次数,和原有阻塞逻辑语义对齐 .take(maxAttempts) // 忽略中间过程的失败列表,直接返回完成信号 .then(); }
若需要获取最终失败的条目,调整实现即可:
public Mono<List<String>> sendWithRetriesGetFinalFailed(List<String> initialItems, int maxAttempts) { return sendsItemsReactive(initialItems) .expand(failedItems -> failedItems.isEmpty() ? Mono.empty() : sendsItemsReactive(failedItems)) .take(maxAttempts) // 取最后一次返回的失败列表,即为经过最大尝试次数后仍发送失败的条目 .last(); }
方案优势:
- 无额外可变状态(比如不需要
AtomicReference),天然线程安全 - 不需要抛出异常触发重试,避免不必要的异常栈开销,逻辑清晰易懂
- 可扩展性强,很容易添加重试间隔、动态重试次数、失败指标上报等附加逻辑
- 完全适配Kinesis PutRecords场景,直接将接口返回的失败记录传入下一次调用即可。
内容的提问来源于stack exchange,提问作者Tomek
相关产品推荐
相关产品推荐

