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

如何无内存分配地从ArrayBlockingQueue中poll并移除元素

解决ArrayBlockingQueue poll操作的内存分配问题

我在使用ArrayBlockingQueue的poll()方法获取元素时遇到了内存分配异常,核心处理线程代码如下:

public void run() {
    while (true) {
        if (Thread.currentThread().isInterrupted()){
            closeAppender();
            break;
        }

        if (!TASK_BUFFER.isEmpty() && !TASK_ADDRESS.isEmpty() && !TASK_MARKER.isEmpty()){
            buffer = TASK_BUFFER.poll();
            remoteAdd = TASK_ADDRESS.poll();
            marker = TASK_MARKER.poll();

//            if (buffer.hasArray()){  //if byte buffer is used just get array
//                byteArr = buffer.array();
//            }else{
//                byteArr = directArr;  //if direct buffer is used, we get into direct arr
////                Arrays.fill(byteArr, (byte)0);  //reset every value to "null" value
//                buffer.get(byteArr, 0, buffer.remaining());
//            }
//            remaining = buffer.remaining();
//            runTask();
            BUFFER_POOL.returnBuffer(buffer);
        }
    }
    LOGGER.info(mainMarker,"Task Scheduler Thread has terminated...");
}

通过jvisualvm监控发现,每次调用poll()处理任务时都会产生内存分配。经过反复注释测试确认,只有poll()和调度操作会影响内存分配量。我的场景不需要遍历整个队列,仅需每次取出一个元素处理并移除,想找到正确的实现方式。

补充信息:

  • 存在单独的生产线程,负责创建ByteBuffer、InetSocketAddress和Marker实例,并分别放入对应的ArrayBlockingQueue
  • 查看了poll()和dequeue()的源码,未发现显式触发内存分配的逻辑。按理论,对象均由生产线程创建并入队,处理线程仅做取出处理,不应产生新的内存分配

队列初始化及生产方法代码如下:

private final BlockingQueue<ByteBuffer> TASK_BUFFER;
private final BlockingQueue<InetSocketAddress> TASK_ADDRESS;
private final BlockingQueue<Marker> TASK_MARKER;
private final BufferPool BUFFER_POOL;
private final ExcerptAppender APPENDER;

public TaskScheduler(BufferPool bufferPool, ChronicleQueue chronicleQueue){
    TASK_BUFFER = new ArrayBlockingQueue<>(bufferPool.getCapacity());
    TASK_ADDRESS = new ArrayBlockingQueue<>(bufferPool.getCapacity());
    TASK_MARKER = new ArrayBlockingQueue<>(bufferPool.getCapacity());
    BUFFER_POOL = bufferPool;
    APPENDER = chronicleQueue.acquireAppender();
}
public void scheduleTask(final ByteBuffer buffer, final InetSocketAddress remoteAdd, final Marker marker) throws InterruptedException {
    TASK_BUFFER.put(buffer);
    TASK_ADDRESS.put(remoteAdd);
    TASK_MARKER.put(marker);
}

问题分析与优化方案

1. 核心问题:多队列同步不一致导致空轮询

当前代码通过!isEmpty()判断三个队列是否有元素,再调用poll(),但这个判断和poll()操作不是原子性的。可能出现判断后某个队列变为空,导致poll()返回null,后续的BUFFER_POOL.returnBuffer(buffer)等操作可能因处理null产生隐性内存分配;同时空轮询会让线程持续自旋,触发JVM的循环优化或锁操作产生临时对象。

2. 最优优化:合并多队列为单队列,封装任务对象

将三个绑定的元素封装成一个Task对象,用单个队列管理,从根源上避免多队列同步问题,减少锁操作次数:

// 封装任务对象
static class Task {
    private final ByteBuffer buffer;
    private final InetSocketAddress remoteAdd;
    private final Marker marker;

    public Task(ByteBuffer buffer, InetSocketAddress remoteAdd, Marker marker) {
        this.buffer = buffer;
        this.remoteAdd = remoteAdd;
        this.marker = marker;
    }

    // 提供getter方法
    public ByteBuffer getBuffer() { return buffer; }
    public InetSocketAddress getRemoteAdd() { return remoteAdd; }
    public Marker getMarker() { return marker; }
}

修改队列初始化与生产方法:

private final BlockingQueue<Task> TASK_QUEUE;

public TaskScheduler(BufferPool bufferPool, ChronicleQueue chronicleQueue){
    TASK_QUEUE = new ArrayBlockingQueue<>(bufferPool.getCapacity());
    BUFFER_POOL = bufferPool;
    APPENDER = chronicleQueue.acquireAppender();
}

public void scheduleTask(final ByteBuffer buffer, final InetSocketAddress remoteAdd, final Marker marker) throws InterruptedException {
    TASK_QUEUE.put(new Task(buffer, remoteAdd, marker));
}

3. 替换isEmpty()+poll()为阻塞式take()

使用take()方法阻塞等待队列元素,避免空轮询的CPU消耗与隐性内存分配,同时保证线程在无任务时进入等待状态:

public void run() {
    try {
        while (!Thread.currentThread().isInterrupted()) {
            // 阻塞等待任务,无需空轮询
            Task task = TASK_QUEUE.take();
            ByteBuffer buffer = task.getBuffer();
            InetSocketAddress remoteAdd = task.getRemoteAdd();
            Marker marker = task.getMarker();

            // 执行任务逻辑
            // runTask();
            BUFFER_POOL.returnBuffer(buffer);
        }
    } catch (InterruptedException e) {
        // 恢复中断状态,保证线程正确终止
        Thread.currentThread().interrupt();
    } finally {
        closeAppender();
        LOGGER.info(mainMarker,"Task Scheduler Thread has terminated...");
    }
}

4. 排查隐性内存分配点

  • 检查BUFFER_POOL.returnBuffer(buffer)方法内部是否存在内存分配逻辑,比如创建临时对象、集合操作等
  • 若仍有内存分配,可通过JVM参数-XX:+PrintGC或专业工具(如AsyncProfiler)定位具体分配点,确认是否为JIT编译或锁操作产生的临时对象

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

相关产品推荐
方舟 Agent Plan

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

最近更新时间:2026.07.29 08:10:09