Flux为空时如何正确使用Consumer/Runnable实现日志记录?
如何在Flux为空时记录日志
你遇到的问题是:尝试通过switchIfEmpty抛出异常,再用onErrorContinue捕获来记录空流日志,但异常未被自定义逻辑处理,反而触发了Reactor的默认错误处理器。
问题原因
onErrorContinue的设计目标是处理上游操作符在处理元素过程中抛出的异常,它允许流跳过出错的元素继续处理后续内容。但空流没有任何元素需要处理,此时抛出的异常属于流的终止信号,onErrorContinue不会捕获这类场景的异常。此外,通过抛出异常再捕获的方式来触发日志,本身就是冗余且不符合Reactor设计理念的做法。
正确实现方式
方案一:直接在switchIfEmpty中执行日志(推荐)
不需要借助异常,直接在空流分支里执行日志逻辑,这是最简洁高效的方式:
import lombok.SneakyThrows; import lombok.extern.slf4j.Slf4j; import org.junit.jupiter.api.Test; import reactor.core.publisher.Flux; import reactor.core.publisher.Mono; @Slf4j public class GenericTest { @Test @SneakyThrows void test() { Flux.empty() .switchIfEmpty(Mono.fromRunnable(() -> log.warn("Flux is empty"))) .subscribe(); Thread.sleep(5000); } }
Mono.fromRunnable会执行日志记录逻辑,同时它本身是一个不发射元素的空Mono,不会改变原流的空状态,完美适配空流场景的日志需求。
方案二:通过错误处理逻辑捕获异常(不推荐,仅作演示)
如果一定要保留异常触发的逻辑,可以使用onErrorResume来捕获终止异常并处理:
import lombok.SneakyThrows; import lombok.extern.slf4j.Slf4j; import org.junit.jupiter.api.Test; import reactor.core.publisher.Flux; import reactor.core.publisher.Mono; @Slf4j public class GenericTest { @Test @SneakyThrows void test() { Flux.empty() .switchIfEmpty(Mono.error(new IllegalStateException())) .onErrorResume(e -> { log.warn("Flux is empty"); return Flux.empty(); }) .subscribe(); Thread.sleep(5000); } }
onErrorResume会捕获流中的终止异常,执行日志逻辑后返回一个新的空流,避免异常扩散到Reactor的默认错误处理器。
内容的提问来源于stack exchange,提问作者Sergey Zolotarev
相关产品推荐
相关产品推荐

