LinkedBlockingQueue为空时生产线程仍阻塞的问题排查求助
问题:LinkedBlockingQueue生产者线程莫名阻塞
为排查问题,我只设置了1个生产线程和1个消费线程,二者共享容量为8的LinkedBlockingQueue。消费线程从队列取到实体后,交给其他线程存入数据库。但经过几次插入、取数迭代后,生产线程突然停止工作。调试日志显示,哪怕队列是空的或者有足够剩余空间,queue.put(entity)操作依然会被阻塞。
生产者代码
final BlockingQueue<Entity> queue = new LinkedBlockingQueue<>(8); //located in calling method // 省略其他代码 do { List<Entity> entityList = entityDatasource.getEntity(); for (Entity entity: entityList) { try { log.debug("Size before insert opertaion is: " + queue.size()); queue.put(entity); log.debug("Size after insert opertaion is: " + queue.size()); } catch (InterruptedException ex) { // 省略异常处理 } } } while (atomicBool.get());
消费者代码
CompletableFuture<Void> queueHandler = CompletableFuture.runAsync(() -> { do { try { log.debug("Queue size is: " + queue.size()); Entity entity = queue.take(); log.debug("Queue size is: " + queue.size()); storeInDb(entity); } catch (InterruptedException ex) { // 省略异常处理 } } while (atomicBool.get()); }, asyncPoolQueueHandler); //ThreadPoolTaskExecutor List<CompletableFuture<Void>> pool = new ArrayList<>(); IntStream.range(0, 1).forEach(i -> { pool.add(queueHandler); }); CompletableFuture.allOf(pool.toArray(CompletableFuture[]::new));
DB存储代码
CompletableFuture .supplyAsync(() -> { return entityRep.save(entity); }, asyncPoolDbPerformer).join(); //ThreadPoolTaskExecutor
我用VisualVM查看了线程状态,没发现明显异常,但生产线程阻塞后,整个处理流程完全停滞。
排查建议
- 检查
atomicBool的状态:确认生产线程和消费线程的终止条件atomicBool.get()是否被意外置为false?如果生产线程的循环条件提前不满足,会直接退出循环,看起来像是“停止工作”,但这是正常终止而非阻塞。建议在循环结束处添加日志,确认是否因atomicBool状态变化导致循环退出。 - 排查
storeInDb的join()阻塞:storeInDb中调用了CompletableFuture.join(),如果asyncPoolDbPerformer线程池耗尽(比如所有线程都卡在DB操作上),supplyAsync无法获取线程执行任务,join()会一直阻塞消费线程。消费线程被堵后无法继续从队列取数据,队列很快会被填满,最终导致生产者的put()阻塞。建议给storeInDb添加日志,记录每个DB操作的开始和结束时间,排查是否有长时间未完成的DB请求;同时检查asyncPoolDbPerformer的线程池配置(核心线程数、最大线程数、队列容量),确认是否因线程池饱和导致任务无法执行。 - 检查
entityDatasource.getEntity()是否阻塞:如果数据源获取数据的方法本身长时间阻塞,也会让生产线程看起来停止工作。建议在调用getEntity()前后添加日志,统计该方法的耗时。 - 完善线程中断处理:代码中捕获了
InterruptedException但未做有效处理,建议在捕获异常后重新设置线程的中断状态(Thread.currentThread().interrupt()),并记录详细的中断日志,避免因线程被中断却未正确处理导致的异常流程。 - 替换
queue.size()的调试方式:LinkedBlockingQueue.size()是弱一致性方法,多线程环境下返回的大小可能不准确。可以改用queue.remainingCapacity()查看剩余容量,或者通过日志记录put和take的调用次数,更准确地跟踪队列实际状态。
内容的提问来源于stack exchange,提问作者ordinary user
相关产品推荐
相关产品推荐

