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

SpringBoot多线程队列消费异常:运行一段时间后仅单线程工作

问题描述

运行SpringBoot应用时,业务服务每秒触发15-30次onMessage方法往全局队列存消息,同时检查全局布尔数组threads,当0-7索引对应值为false时启动对应线程。每个线程从队列取消息处理,运行1分钟后将threads对应索引设为false。

应用初期8个线程均正常输出处理日志,但运行5-10分钟后,仅单个线程持续输出,其余线程停滞。尝试过线程休眠、调整continue逻辑,未找到停滞原因,怀疑同步问题但无法确定,寻求排查指导。

相关代码片段

消息入队与线程触发代码

public static void onMessage(String record) {
     global.add(record);
     if(threads[0] == false) {
     threads[0] = true;
     thread0.start() // Name of the thread index included in runnable
}

线程run方法代码

public void run() { 
    String recordToUse;
    int thread_num = Integer.parseInt(Thread.currentThread().getName());
    long startThreadTime = System.currentTimeMillis();
    long endThreadTime = startThreadTime + 60 * 1000; // Run the thread for 1 minute.
                
    while(System.currentTimeMillis() < endThreadTime) {
        if(!global.isEmpty()) {
            recordToUse = global.remove();
            System.out.println("Successful removal: Thread-"+ thread_num);
        } else {
            continue; // If the queue is empty, keep checking until it is not empty. 
        }
        // 消息处理操作
    }
    threads[thread_num] = false;
    return;
}

线程实例化示例:Thread thread1 = new Thread(runnable,"0")

排查与解决方案

1. 核心问题:非线程安全的全局变量

  • 队列非线程安全:如果global是普通ArrayList这类非线程安全集合,多线程并发add()和remove()会破坏队列内部结构,比如导致remove()抛出异常、数据丢失,线程可能因未捕获异常静默终止。
  • 布尔数组无同步保护:threads数组的读写没有同步机制,多个线程操作时会出现可见性问题——线程修改threads的值后,其他线程无法及时感知,错误判断线程状态,比如误以为线程未启动重复调用start()(线程只能启动一次,重复调用会抛IllegalThreadStateException终止线程)。

修复方案:

  • 替换global为线程安全队列,比如LinkedBlockingQueue,自带阻塞等待逻辑,无需手动空轮询。
  • 替换threads为AtomicBoolean[]原子数组,用compareAndSet()保证状态判断和更新的原子性,确保多线程下状态可见。

2. 空轮询引发的线程饥饿

线程在队列空时用continue无限空轮询,会把CPU核心占满,导致其他线程得不到调度机会,时间长了就会出现只有单个线程能抢到CPU运行的情况。

修复方案:

  • 用线程安全队列的take()方法替代空轮询,take()会在队列空时自动阻塞线程,直到有新消息进来,既节省CPU,又能保证线程及时被唤醒处理消息。

3. 未捕获异常导致线程静默死亡

如果消息处理逻辑中抛出未捕获的异常,线程会直接终止,但代码中没有异常捕获,线程死了也不会输出日志,看起来就像停滞了。

修复方案:

  • 在run()方法中给消息处理逻辑加try-catch块,捕获所有异常并打印日志,排查是否是业务逻辑抛异常导致线程终止。

调整后的示例代码

全局变量定义

// 线程安全阻塞队列
private static final BlockingQueue<String> global = new LinkedBlockingQueue<>();
// 原子布尔数组,保证线程安全的状态读写
private static final AtomicBoolean[] threads = new AtomicBoolean[]{
    new AtomicBoolean(false),
    new AtomicBoolean(false),
    new AtomicBoolean(false),
    new AtomicBoolean(false),
    new AtomicBoolean(false),
    new AtomicBoolean(false),
    new AtomicBoolean(false),
    new AtomicBoolean(false)
};

消息入队与线程触发逻辑

public static void onMessage(String record) {
    global.add(record);
    // 以线程0为例,用compareAndSet原子操作判断并更新状态,避免重复启动
    if (threads[0].compareAndSet(false, true)) {
        new Thread(new MessageRunnable(), "0").start();
    }
    // 其他线程(1-7)的触发逻辑同理,替换对应的索引即可
}

线程run方法优化

public class MessageRunnable implements Runnable {
    @Override
    public void run() {
        int threadNum = Integer.parseInt(Thread.currentThread().getName());
        long endThreadTime = System.currentTimeMillis() + 60 * 1000;

        while (System.currentTimeMillis() < endThreadTime) {
            try {
                // 阻塞等待消息,替代空轮询
                String recordToUse = global.take();
                System.out.println("Successful removal: Thread-" + threadNum);
                // 消息处理逻辑
                processMessage(recordToUse);
            } catch (InterruptedException e) {
                // 线程被中断时重置状态并退出
                Thread.currentThread().interrupt();
                break;
            } catch (Exception e) {
                // 捕获业务异常,避免线程静默死亡
                System.err.println("Thread-" + threadNum + "处理消息失败: " + e.getMessage());
                e.printStackTrace();
            }
        }
        // 线程结束后更新状态
        threads[threadNum].set(false);
    }

    private void processMessage(String record) {
        // 你的业务处理逻辑
    }
}

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

相关产品推荐
方舟 Agent Plan

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

最近更新时间:2026.08.13 05:10:32