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
相关产品推荐
相关产品推荐

