Spring WebFlux中map内调用Repository执行block永久挂起问题
问题描述
响应式编程实现用户对象映射转换时出现程序永久挂起问题:
- 预期逻辑:根据用户名到数据库查询匹配用户,创建新用户对象,调用另一个Repository查询权限信息填充新用户的Authority属性,最终调用
block()方法获取返回值。 - 故障表现:执行到
block()方法时程序卡住,数据库请求始终无法完成。调试发现进入Mono#block下的BlockingSingleSubscriber#blockingGet逻辑时,getCount()返回值为1,进入await()无限等待,计数始终无法降至0,也不会抛出InterruptedException。硬编码角色值(不访问数据库查询)时程序可正常运行。
故障相关代码
用户对象转换逻辑:
Mono<User> userMono = userService.findByUsername(regularUser.getUsername()) .map(u -> { User user = new User(); user.setUsername(u.getUsername()); user.setAuthority(userService.getAuthorities(u)); return user; }); User usr = userMono.block();
UserService定义:
@Autowired UserRepository userRepository; @Autowired AuthorityRepository authorityRepository; public Mono<User> findByUsername(String username) { return userRepository.findByUsername(username); } public Flux<Authority> getAuthorities(User user) { return authorityRepository.findAllByUserId(user.getId()); }
存在问题的ReactiveUserDetailsService定义:
@Bean public PasswordEncoder passwordEncoder() { return PasswordEncoderFactories.createDelegatingPasswordEncoder(); } @Bean public ReactiveAuthenticationManager reactiveAuthenticationManager(ReactiveUserDetailsService userDetailsService, PasswordEncoder passwordEncoder) { var authenticationManager = new UserDetailsRepositoryReactiveAuthenticationManager(userDetailsService); authenticationManager.setPasswordEncoder(passwordEncoder); return authenticationManager; } @Bean public ReactiveUserDetailsService userDetailsService(UserService userService) { return username -> userService.findByUsername(username) .log() .map(u -> org.springframework.security.core.userdetails.User .withUsername(u.getUsername()).password(u.getPassword()) .roles(userService.getRoles(u).toArray(String[]::new)) .accountExpired(!u.isActive()) .credentialsExpired(!u.isActive()) .disabled(!u.isActive()) .accountLocked(!u.isActive()) .build() ); }
使用技术栈:Spring Boot 2.7.0、io.spring.dependency-management 1.0.11.RELEASE、spring-boot-starter-webflux、spring-boot-starter-data-mongodb-reactive。
故障根因
两个核心错误导致死锁:
- 响应式算子使用错误:
map算子仅用于同步转换数据,在map内部直接调用返回Flux/Mono的响应式Repository方法,既没有订阅嵌套的响应式流,也没有通过专门的异步衔接算子将嵌套流纳入主链路,嵌套的数据库查询根本不会被执行。 - 阻塞操作触发线程死锁:Spring WebFlux + Reactive MongoDB默认使用固定数量的Netty事件循环线程调度所有响应式任务,在事件循环线程上直接调用
block()会直接占住线程等待结果,而数据库查询的响应任务需要事件循环线程才能执行,线程被占住后任务永远无法调度,就会出现无限等待的情况。ReactiveUserDetailsService里直接调用userService.getRoles(u).toArray()本质也是对响应式流做隐式阻塞,同样会触发线程卡死。
修复方案
- 算子替换:所有需要基于前序流结果发起新的响应式调用(返回Mono/Flux)的场景,用
flatMap/flatMapMany替代map,将嵌套的异步流纳入主响应式链路,保证流会被正确订阅执行。 - 阻塞操作限制:禁止在WebFlux请求处理、响应式数据访问的链路内部调用
block(),block()仅能用于响应式代码和非响应式代码的边界场景(比如应用启动初始化、单元测试、命令行任务),避免抢占事件循环线程。 - 集合类型结果处理:Flux返回的多值结果要通过
collectList()等算子在流链路内完成收集,不要直接在map里阻塞拉取。
修复后的用户对象转换代码:
Mono<User> userMono = userService.findByUsername(regularUser.getUsername()) .flatMap(u -> userService.getAuthorities(u) .collectList() .map(authorities -> { User user = new User(); user.setUsername(u.getUsername()); // 这里传入的是实际查询到的权限集合,不是未订阅的Flux对象 user.setAuthority(authorities); return user; }) ); // 仅在启动初始化、测试等非请求线程场景使用block() User usr = userMono.block();
修复后的ReactiveUserDetailsService代码:
@Bean public ReactiveUserDetailsService userDetailsService(UserService userService) { return username -> userService.findByUsername(username) .log() .flatMap(u -> userService.getRoles(u) .collectList() .map(roles -> org.springframework.security.core.userdetails.User .withUsername(u.getUsername()) .password(u.getPassword()) .roles(roles.toArray(String[]::new)) .accountExpired(!u.isActive()) .credentialsExpired(!u.isActive()) .disabled(!u.isActive()) .accountLocked(!u.isActive()) .build() ) ); }
注意:如果是在WebFlux的接口处理逻辑中,直接把组装完成的
Mono/Flux返回给框架即可,框架会自动完成订阅和结果返回,不需要手动调用block()。
内容的提问来源于stack exchange,提问作者user3807806
相关产品推荐
相关产品推荐

