WebFlux AOP场景下订阅Mono无输出,如何不返回Publisher即可订阅?
问题根因
你提供的代码没有日志输出,核心原因有两点:
ReactiveSecurityContextHolder的上下文是绑定在原请求的响应式流订阅链上的,你在AOP通知中单独订阅新流时没有继承原请求上下文,无法拿到安全上下文,流执行时要么返回空值要么抛出异常,你没有配置错误回调所以无法感知到问题,自然没有日志输出。- 如果你切点匹配的是返回
Mono/Flux的WebFlux接口方法,@AfterReturning执行时原业务流还没有被订阅执行,你单独启动的订阅流和原业务流完全隔离,拿不到运行时的上下文数据。
解决方案
方案1:环绕通知绑定侧写逻辑(推荐)
如果你的切点方法返回值是Mono/Flux类型,优先选择该方案,不需要手动管理订阅,侧写逻辑跟随原业务流生命周期执行,不会丢失上下文:
@Aspect @Component public class LogAop { // 匹配返回Mono/Flux的业务方法 @Pointcut("execution(public reactor.core.publisher.Mono com.yourpackage..*(..)) || execution(public reactor.core.publisher.Flux com.yourpackage..*(..))") public void reactiveServicePointcut() {} @Around("reactiveServicePointcut()") public Object aroundReactiveMethod(ProceedingJoinPoint pjp) throws Throwable { Object result = pjp.proceed(); if (result instanceof Mono<?>) { return ((Mono<?>) result) // 业务流正常完成后执行侧写 .doOnSuccess(data -> { // 直接在原流上下文中获取用户名,不需要单独传递上下文 ReactiveSecurityContextHolder.getContext() .map(ctx -> ctx.getAuthentication().getName()) .doOnNext(username -> { log.info("操作用户名:{}", username); // 存库操作直接订阅,添加错误回调避免影响主业务 saveOperationLog(username) .onErrorComplete(err -> { log.error("存操作日志失败", err); return true; }) .subscribe(); }) .onErrorComplete() .subscribe(); }) // 业务流报错也可以记录日志 .doOnError(err -> log.error("业务执行失败", err)); } else if (result instanceof Flux<?>) { return ((Flux<?>) result) .doOnComplete(() -> { // Flux场景的日志/存库逻辑和Mono一致 }) .doOnError(err -> log.error("业务执行失败", err)); } return result; } // 示例存库方法 private Mono<Void> saveOperationLog(String username) { OperationLog log = OperationLog.builder() .username(username) .operateTime(LocalDateTime.now()) .build(); return reactiveMongoTemplate.save(log).then(); } }
方案2:独立订阅(无需返回流给上层)
如果确实必须在@AfterReturning中执行逻辑,不需要把流返回给上层,可以手动传递上下文、配置调度器和异常回调:
@AfterReturning("execution(* com.yourpackage..*(..))") public void configSetted() { ReactiveSecurityContextHolder.getContext() .map(ctx -> ctx.getAuthentication().getName()) // 用弹性线程池调度IO操作,避免占用核心IO线程 .publishOn(Schedulers.boundedElastic()) .subscribe( username -> { log.info("操作用户名:{}", username); // 执行存库操作 saveOperationLog(username) .subscribe(null, err -> log.error("存操作日志失败", err)); }, // 添加上下文获取的错误回调,方便排查问题 err -> log.error("获取用户名失败", err) ); }
注意事项
- 所有非核心业务的侧写操作(日志、存库)必须添加异常处理逻辑,避免报错影响主业务流程。
- IO密集型的侧写操作需要指定
Schedulers.boundedElastic()调度,避免占用WebFlux的核心IO线程。 - 不要在响应式流的执行链中调用
block()方法阻塞线程,否则会引发线程耗尽问题。
内容的提问来源于stack exchange,提问作者Morteza Malvandi
相关产品推荐
相关产品推荐

