Reactor中带退避的RetryWhen无效问题排查求助
问题:Reactor带退避的延迟重试未生效,无重试日志输出
尝试实现带退避的延迟重试功能,复制官方RetryWhen示例代码后,运行发现重试未生效,日志里完全没有doAfterRetry相关输出。
代码
AtomicInteger errorCount = new AtomicInteger(); Flux<String> flux = Flux.<String>error(new IllegalStateException("boom")) .doOnError(e -> { errorCount.incrementAndGet(); System.out.println(e + " at " + LocalTime.now()); }) .retryWhen(Retry .backoff(3, Duration.ofMillis(100)) .jitter(0d) .doAfterRetry(rs -> System.out.println("retried at " + LocalTime.now() + ", attempt " + rs.totalRetries())) .onRetryExhaustedThrow((spec, rs) -> rs.failure()) ); flux.subscribe();
运行日志
2023-01-10 18:54:43.851 [main] DEBUG reactor.util.Loggers.debug-254 - Using Slf4j logging framework java.lang.IllegalStateException: boom at 18:54:43.873 Process finished with exit code 0
环境信息
java version "1.8.0_281" Java(TM) SE Runtime Environment (build 1.8.0_281-b09) Java HotSpot(TM) 64-Bit Server VM (build 25.281-b09, mixed mode) reactor-core : 3.5.0
解答
核心问题是Reactor的异步非阻塞特性:flux.subscribe()是异步调用,主线程执行完这个方法后直接退出,此时退避延迟还没结束,重试逻辑根本没机会执行,程序就终止了。
解决办法
让主线程等待异步操作完成即可,以下是几种常用方案:
简单测试场景用
block()
修改订阅代码,用blockLast()阻塞主线程直到Flux执行完成:flux.blockLast();手动同步用
CountDownLatch
通过计数锁让主线程等待重试流程结束:CountDownLatch latch = new CountDownLatch(1); flux.subscribe( s -> {}, e -> latch.countDown(), latch::countDown ); latch.await();单元测试用
StepVerifier
适合测试环境下验证重试逻辑:StepVerifier.create(flux) .expectError(IllegalStateException.class) .verify();
修改后重新运行,就能看到doAfterRetry的日志输出,重试逻辑也会按预期执行。
内容的提问来源于stack exchange,提问作者Peng
相关产品推荐
相关产品推荐

