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

Spring WebFlux Mono实现两次过滤 二次过滤不执行问题排查

问题根因

第二层filter不执行的核心原因是响应式流的空值传播特性:

  • Spring Data Reactive的findByEmail()方法在查询不到对应邮箱的用户时,返回的是空Mono(Mono.empty()),而不是包装了null的Mono。
  • zipWith操作符的规则是:只要参与组合的任意一个Mono为空,组合后的结果Mono就会直接为空,不会向下游发射任何元素。你当前代码里当邮箱不存在(也就是你期望走注册逻辑的场景)时,findByEmail返回空Mono,zip后的Tuple流直接结束,后续的filter拿不到任何待处理元素,自然不会执行。

反过来如果邮箱已经存在,findByEmail能查到用户,zip才会发射元素,这时候filter判断objects.getT1() != null就会把流过滤空,后面的保存逻辑永远不会触发——刚好和你预期的逻辑反过来了。

修复方案

给findByEmail的返回结果加空值兜底,把查询不到用户的空Mono转换为发射null值的Mono,保证不管查没查到用户,zip操作都能拿到元素向下传递,filter才能正常执行判断。

修正后可运行代码
public Mono<String> signupUser(Mono<UserTransfer> userMono) {
        LOG.info("signup user");

        return userMono
                .filter(userTransfer -> {
                    // 第一层过滤:校验apiKey匹配
                    LOG.info(" userTransfer.apiKey: {}, apiKey: {}, match: {}", userTransfer.getApiKey(), apiKey, userTransfer.getApiKey().equals(apiKey));
                    return userTransfer.getApiKey().equals(apiKey);
                })
                // 可选优化:apiKey不匹配时返回明确错误,避免前端拿到无信息的404响应
                // .switchIfEmpty(Mono.error(() -> new ResponseStatusException(HttpStatus.UNAUTHORIZED, "invalid api key")))
                .flatMap(userTransfer -> {
                    // 关键修改:查询无结果时兜底返回null,避免空Mono中断zip流程
                    Mono<User> existingUserMono = userRepository.findByEmail(userTransfer.getEmail())
                            .defaultIfEmpty(null);
                    
                    LOG.info("do findByEmail for email: {}", userTransfer.getEmail());
                    return existingUserMono.zipWith(Mono.just(userTransfer));
                })              
                .filter(objects -> {
                    // 第二层过滤:校验邮箱未被注册
                    LOG.info("user findByEmail result is {}, null means email is available for signup", objects.getT1());
                    return objects.getT1() == null;
                } )
                // 可选优化:邮箱已存在时返回明确错误
                // .switchIfEmpty(Mono.error(() -> new ResponseStatusException(HttpStatus.CONFLICT, "email already registered")))
                .flatMap(objects -> {
                    User user = new User(objects.getT2().getFirstName(), objects.getT2().getLastName(), objects.getT2().getEmail());
                    return userRepository.save(user).zipWith(Mono.just(objects.getT2()));
                    })
                .flatMap(objects -> {
                    User user = objects.getT1();
                    UserTransfer userTransfer = objects.getT2();
                    LOG.info("create authentication with rest call to endpoint {}", authenticationEp);

                    WebClient.ResponseSpec responseSpec = webClient.post().uri(authenticationEp).bodyValue(userTransfer).retrieve();
                    return responseSpec.bodyToMono(String.class).map(authenticationId -> {
                        LOG.info("got back authenticationId from service call: {}", authenticationId);
                        return "user created with id: " + user.getEmail() + " with authId: " + authenticationId;
                    });
        });
    }
额外说明

不建议在响应式流中随意传递null值,更规范的实现是用Optional<User>包装查询结果,通过objects.getT1().isEmpty()判断邮箱是否可用,从根源上避免空指针风险。

内容的提问来源于stack exchange,提问作者Katlock

相关产品推荐
方舟 Agent Plan

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

最近更新时间:2026.09.02 07:12:52