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

ParallelStream结合JPA查询返回空结果问题排查

ParallelStream 执行时数据库查询结果不符合预期的原因分析

问题场景

我编写了如下单元测试和业务代码:

单元测试类

@DataJpaTest
@TestExecutionListeners({ TransactionalTestExecutionListener.class, DependencyInjectionTestExecutionListener.class,
        DbUnitTestExecutionListener.class, MockitoTestExecutionListener.class })
@DatabaseSetup(value = "/dataset/myDataSet.xml", type = DatabaseOperation.CLEAN_INSERT)
@DatabaseTearDown(value = "/dataset/myDataSet.xml", type = DatabaseOperation.DELETE_ALL)
public class MyTest {
    private ClassToBeTested classToTest;
    
    @MockBean
    private MyObjRepository myObjRepo;
    @MockBean
    private MyRepository myRepo;
    
    @BeforeEach
    public void setup() {
        MockitoAnnotations.initMocks(this);
        classToTest = new ClassToBeTested(myRepo);
    }
        
    @Test
    void myTest() {
        //数据库中已添加2个对象
        classToTest.scheduleFixedDelayTask();
    }
}

Repository 接口

public interface MyRepository extends JpaRepository<OtherObject, Long> {
}

业务服务类

@Service
@RequiredArgsConstructor
@Slf4j
public class ClassToBeTested {
    private final MyObjRepository myObjRepo;
    private final MyRepository myRepo;
    
    
    @Scheduled(fixedDelayString = "${my.var}")
    public void scheduleFixedDelayTask() {
        manage();
    }

    protected void manage() {
        MyObj obj1 = myObjRepo.findByName("obj1");
        MyObj obj2 = myObjRepo.findByName("obj2");
        Map<String, List<MyObj>> myMap = new ConcurrentHashMap<>();
        myMap.put("obj1", Collections.singletonList(obj1));
        myMap.put("obj2", Collections.singletonList(obj2));
        myMap.entrySet().parallelStream().forEach(items -> {
            for (MyObj item : items.getValue()) {
                List<OtherObject> otherObjects = MyRepository.findAll();
                log.info("TEST: " + CollectionUtils.isEmpty(otherObjects));
            }
        });
    }
}

现象

  • 使用parallelStream时,日志输出:
    TEST: true
    TEST: false
    
    不符合预期(预期应为false/false)
  • 使用普通stream时,日志输出正常:
    TEST: false
    TEST: false
    
    注:obj1和obj2均不为空

原因分析

核心问题:静态调用仓库方法 + 线程上下文隔离

问题出在业务代码的这一行:

List<OtherObject> otherObjects = MyRepository.findAll();
  1. Spring Data JPA 仓库的线程绑定特性:Spring托管的Repository实例依赖于ThreadLocal管理的EntityManager上下文,这个上下文只在当前请求/测试线程中有效。
  2. 普通stream的执行逻辑:在当前测试线程中执行,能正确获取到Spring注入的myRepo实例及其绑定的EntityManager,因此可以查询到数据库中预先填充的测试数据,返回非空集合。
  3. parallelStream的执行逻辑:会使用线程池中的多线程并行处理任务,这些线程没有被Spring初始化,无法获取到测试线程的EntityManager上下文。直接静态调用MyRepository.findAll()会创建一个未绑定测试数据的仓库实例,导致查询不到数据,返回空集合。

另外还有一个次要问题:测试类中ClassToBeTested的构造方法调用存在参数缺失(原类依赖两个Repository,但构造时只传了myRepo),不过这不是当前现象的直接原因,因为myObjRepo已被Mock且obj1、obj2不为空。

解决方案

  1. 替换静态调用为实例方法:使用类中注入的myRepo实例调用findAll(),这样就能在并行线程中正确使用Spring托管的仓库实例(Spring Data JPA的Repository默认是线程安全的):
    List<OtherObject> otherObjects = myRepo.findAll();
    
  2. 若必须在并行流中使用Spring Bean,确保上下文传递:可以通过RequestContextHolder手动复制上下文到并行线程,但这种方式较为繁琐,优先推荐使用实例方法调用的方式。

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

相关产品推荐
方舟 Agent Plan

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

最近更新时间:2026.07.07 10:57:46