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

