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

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

相关产品推荐
方舟 Agent Plan

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

最近更新时间:2026.08.06 10:25:13