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

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。

故障根因

两个核心错误导致死锁:

  1. 响应式算子使用错误:map算子仅用于同步转换数据,在map内部直接调用返回Flux/Mono的响应式Repository方法,既没有订阅嵌套的响应式流,也没有通过专门的异步衔接算子将嵌套流纳入主链路,嵌套的数据库查询根本不会被执行。
  2. 阻塞操作触发线程死锁:Spring WebFlux + Reactive MongoDB默认使用固定数量的Netty事件循环线程调度所有响应式任务,在事件循环线程上直接调用block()会直接占住线程等待结果,而数据库查询的响应任务需要事件循环线程才能执行,线程被占住后任务永远无法调度,就会出现无限等待的情况。ReactiveUserDetailsService里直接调用userService.getRoles(u).toArray()本质也是对响应式流做隐式阻塞,同样会触发线程卡死。
修复方案
  1. 算子替换:所有需要基于前序流结果发起新的响应式调用(返回Mono/Flux)的场景,用flatMap/flatMapMany替代map,将嵌套的异步流纳入主响应式链路,保证流会被正确订阅执行。
  2. 阻塞操作限制:禁止在WebFlux请求处理、响应式数据访问的链路内部调用block(),block()仅能用于响应式代码和非响应式代码的边界场景(比如应用启动初始化、单元测试、命令行任务),避免抢占事件循环线程。
  3. 集合类型结果处理: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

相关产品推荐
方舟 Agent Plan

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

最近更新时间:2026.09.03 10:54:32