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
相关产品推荐
相关产品推荐

