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

ExecutorService线程批量拉取数据时查询结果异常问题

多线程批量拉取数据时查询结果条数异常的问题排查与解决

首先,我一眼就注意到你遇到的核心问题:同一个firstResult值,单独调用能返回预期的1000条数据,但在多线程循环提交任务时,返回条数却少于预期,而且这个问题出现在多个firstResult值上。结合你贴的代码,问题的根源其实很明确——JPA的TypedQuery实例不是线程安全的!

问题原因分析

看你的ExplicitAudienceCreationExecutorService类,fetchFromMartAndInsertIntoMongo方法接收一个TypedQuery<Object[]> qr参数,然后在每个线程的Runnable里直接复用这个query对象,调用query.setFirstResult(firstResult).setMaxResults(chunkSize)。

但JPA的Query(包括TypedQuery)实例是有状态的,setFirstResult和setMaxResults都是修改这个实例的内部状态,而且这些方法没有做线程同步。当多个线程同时操作同一个Query实例时,就会出现状态覆盖的情况:

  • 比如线程A刚把firstResult设为302000,还没执行getResultList(),线程B就把同一个query的firstResult改成了303000;
  • 这时候线程A执行查询时,用的是被线程B修改后的参数,自然返回的结果条数就不符合预期了。

解决方案

针对这个问题,最合理的解决方式是保证每个线程使用独立的TypedQuery实例,避免多线程共享同一个有状态的Query对象。下面给你两种可行的修改方案:

方案一:每个线程内创建新的TypedQuery实例

这是最优解,因为TypedQuery是轻量级对象,创建新实例几乎没有性能开销,同时能彻底解决线程安全问题。

你需要修改fetchFromMartAndInsertIntoMongo方法的参数,不再直接传入TypedQuery,而是传入创建Query所需的依赖(比如EntityManager、查询语句、参数等),然后在每个线程的run方法内重新创建Query:

// 修改后的fetchFromMartAndInsertIntoMongo方法示例
public void fetchFromMartAndInsertIntoMongo(int fr, int cs, EntityManager entityManager, String queryString, 
                                            Promotion promotion, FilterKeywords filterKeywords, 
                                            String audienceFilterName, String programId, int queryrernCont) {
    final int firstResult = fr;
    final int chunkSize = cs;

    campaignExecutorService.dotask(new Runnable() {
        @Override
        public void run() {
            mongoTemplate = tenantTemplates.get(programId);
            
            // 关键:每个线程独立创建TypedQuery实例
            TypedQuery<Object[]> query = entityManager.createQuery(queryString, Object[].class);
            // 这里可以根据需要设置查询参数(如果原来的query有参数的话)
            // query.setParameter("paramName", paramValue);
            
            final List<Object[]> toReturn = query.setFirstResult(firstResult).setMaxResults(chunkSize).getResultList();
            
            classCount++;
            System.out.println("classCount "+ classCount);
            logger.info("firstResult "+ firstResult + " queryResultSize " + toReturn.size() );
            
            // 后续的转换和插入MongoDB逻辑不变...
        }
    });
}

方案二:对Query的访问加同步锁(不推荐)

如果因为某些限制,你必须复用同一个TypedQuery实例,可以对query对象的操作加同步锁,保证同一时间只有一个线程修改和执行查询。但这种方式会让所有任务串行执行,完全丧失多线程的并发优势,所以只作为备选方案:

public void fetchFromMartAndInsertIntoMongo(int fr, int cs, TypedQuery<Object[]> qr, 
                                            Promotion promotion, FilterKeywords filterKeywords, 
                                            String audienceFilterName, String programId, int queryrernCont) {
    final int firstResult = fr;
    final int chunkSize = cs;
    final TypedQuery<Object[]> query = qr;

    campaignExecutorService.dotask(new Runnable() {
        @Override
        public void run() {
            mongoTemplate = tenantTemplates.get(programId);
            
            // 加同步锁,同一时间只有一个线程操作query
            synchronized(query) {
                final List<Object[]> toReturn = query.setFirstResult(firstResult).setMaxResults(chunkSize).getResultList();
                
                classCount++;
                System.out.println("classCount "+ classCount);
                logger.info("firstResult "+ firstResult + " queryResultSize " + toReturn.size() );
                
                // 后续逻辑不变...
            }
        }
    });
}

额外检查点

为了彻底确认问题,你可以在日志里多打印一些信息:

  • 打印当前线程ID,看同一个firstResult对应的线程是否有参数被其他线程修改的情况;
  • 确认数据源在查询期间没有数据的增删(不过你单独调用正常,多线程异常,所以这个概率很低)。

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

相关产品推荐
方舟 Agent Plan

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

最近更新时间:2026.05.15 04:02:59