You need to enable JavaScript to run this app.
优惠活动
大模型
产品
解决方案
定价
更多

WebFlux AOP场景下订阅Mono无输出,如何不返回Publisher即可订阅?

问题根因

你提供的代码没有日志输出,核心原因有两点:

  1. ReactiveSecurityContextHolder的上下文是绑定在原请求的响应式流订阅链上的,你在AOP通知中单独订阅新流时没有继承原请求上下文,无法拿到安全上下文,流执行时要么返回空值要么抛出异常,你没有配置错误回调所以无法感知到问题,自然没有日志输出。
  2. 如果你切点匹配的是返回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

相关产品推荐
方舟 Agent Plan

超全模态模型 × Harness 升级,最新支持 Deepseek-V4.1-Flash、GLM-5.3 系列、Doubao-Seedream-5.0-pro、Kimi-K3 (部分), 限时 9.9 元起

最近更新时间:2026.09.27 10:36:02