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

Java中SingleThreadExecutor触发OutOfMemoryError的原因及解决方案

RabbitMQ消息转发服务OOM问题分析与修复

问题场景

有一个RabbitMQ消息生产者,SpringBoot服务作为消费者接收队列消息,将消息存入本地ArrayDeque后,需按接收顺序发送至另一应用的Socket。核心代码如下:

public void addMessageToQueue(CML cml) throws ParseException {

    if (cml != null) {
        AgentEventData agentEventData = setAgentEventData(cml);
        log.info("Populated AgentEventData: {} ", agentEventData);
        MessageProcessor.getMessageQueue().getMessageQueue().add(agentEventData);
        log.info("Message QUEUE Size: {}", MessageProcessor.getMessageQueue().getMessageQueue().size());
        QUEUE_MONITOR.setCachedQueue(MessageProcessor.getMessageQueue());

        executeTasks();

    } else {
        log.error("CML Message is NULL, Message Cannot be added to the Message Queue.");
    }
}

private static void executeTasks() {
    ExecutorService executorService = Executors.newSingleThreadExecutor();
    try {
        executorService.execute(new MessageProcessor());
    } catch (Exception e) {
        log.error("Exception when executing Task: {}", e.getMessage());
    }
    log.info("Shutting down Executor Service........");
    executorService.shutdown();
    log.info("Executor Service Shutdown : {}", executorService.isShutdown());
}

运行一段时间后出现java.lang.OutOfMemoryError: 无法创建本地线程,尝试newFixedThreadPool(10)后问题依旧。

错误原因

  1. 线程池频繁创建销毁:每次调用addMessageToQueue都会新建线程池,执行一次任务就shutdown。线程池shutdown后内部线程不会立即回收,短时间大量消息涌入会导致系统创建的线程数超过操作系统上限,触发OOM。
  2. 任务执行逻辑低效:MessageProcessor仅执行一次就结束,每来一条消息就启动新线程处理,完全失去线程池复用的意义,反而加剧线程资源消耗。
  3. 队列线程不安全:ArrayDeque并非线程安全容器,RabbitMQ消费者多线程处理时,存在并发修改队列的风险。

最优解决方案

1. 复用全局单线程线程池

初始化一个全局唯一的单线程线程池,无需每次创建新池。单线程既保证消息顺序发送,又避免线程资源浪费:

// 全局静态初始化,仅创建一次
private static final ExecutorService MESSAGE_PROCESS_EXECUTOR = Executors.newSingleThreadExecutor();

2. 改造为生产者-消费者模式

将ArrayDeque替换为线程安全的BlockingQueue,让处理线程持续监听队列,有消息时自动处理,避免空轮询:

// 替换原有ArrayDeque为阻塞队列,保证线程安全
private static final BlockingQueue<AgentEventData> MESSAGE_QUEUE = new LinkedBlockingQueue<>();

// MessageProcessor的run方法改为持续监听队列
@Override
public void run() {
    while (!Thread.currentThread().isInterrupted()) {
        try {
            // 阻塞等待队列中有消息
            AgentEventData data = MESSAGE_QUEUE.take();
            // 执行发送到Socket的逻辑
            sendToSocket(data);
        } catch (InterruptedException e) {
            Thread.currentThread().interrupt();
            log.info("消息处理线程被中断");
            break;
        } catch (Exception e) {
            log.error("发送消息到Socket失败", e);
        }
    }
}

// 应用启动时启动一次处理线程即可
@PostConstruct
public void initMessageProcessor() {
    MESSAGE_PROCESS_EXECUTOR.execute(new MessageProcessor());
}

// 修改addMessageToQueue,仅需将消息加入队列
public void addMessageToQueue(CML cml) throws ParseException {
    if (cml != null) {
        AgentEventData agentEventData = setAgentEventData(cml);
        log.info("Populated AgentEventData: {} ", agentEventData);
        MESSAGE_QUEUE.add(agentEventData);
        log.info("Message QUEUE Size: {}", MESSAGE_QUEUE.size());
        QUEUE_MONITOR.setCachedQueue(MESSAGE_QUEUE);
    } else {
        log.error("CML Message is NULL, Message Cannot be added to the Message Queue.");
    }
}

3. 优雅关闭线程池

结合Spring生命周期注解,在应用关闭时优雅终止线程池,避免资源泄漏:

@PreDestroy
public void shutdownExecutor() {
    MESSAGE_PROCESS_EXECUTOR.shutdown();
    try {
        if (!MESSAGE_PROCESS_EXECUTOR.awaitTermination(10, TimeUnit.SECONDS)) {
            MESSAGE_PROCESS_EXECUTOR.shutdownNow();
        }
    } catch (InterruptedException e) {
        MESSAGE_PROCESS_EXECUTOR.shutdownNow();
        Thread.currentThread().interrupt();
    }
}

方案优势

  • 线程资源复用:全程复用一个线程,避免频繁创建销毁线程的资源消耗。
  • 顺序严格保证:单线程处理天然保证消息发送顺序与RabbitMQ接收顺序一致。
  • 线程安全可靠:BlockingQueue保证队列操作的线程安全性,阻塞等待避免空轮询浪费CPU。
  • 资源无泄漏:通过Spring生命周期管理线程池的启停,避免资源残留。

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

相关产品推荐
方舟 Agent Plan

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

最近更新时间:2026.08.12 14:25:45