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

agrona.OneToOneRingBuffer使用异常:大消息量下消费者无输出问题问询

关于Agrona OneToOneRingBuffer的线程安全与生产消费异常问题

我测试Agrona的OneToOneRingBuffer,实现了单生产者单消费者的消息生产消费逻辑,但遇到一个问题:当Config.SIZE设为100_000时,消费者没有输出,CountDownLatch无法归0;设为10_000时程序运行正常。相关代码如下:

public class Producer implements Runnable{
    private final RingBuffer buffer;

    public Producer(RingBuffer buffer) {
        this.buffer = buffer;
    }

    @Override
    public void run() {
        for(int i=0;i<Config.SIZE;i++){
            String s = String.format("i am %s",i);
            System.out.println( "name -> " + s);
            UnsafeBuffer unsafeBuffer = new UnsafeBuffer(s.getBytes());
            unsafeBuffer.wrap(s.getBytes());
            buffer.write(1, unsafeBuffer, 0, s.length());
        }
    }
}

public class Consumer implements Runnable {
    private final RingBuffer buffer;

    public Consumer(RingBuffer buffer) {
        this.buffer = buffer;
    }

    @Override
    public void run() {
        long start = System.currentTimeMillis();
        CountDownLatch countDownLatch = new CountDownLatch(Config.SIZE);
        while (countDownLatch.getCount() > 0) {
            buffer.read((msgTypeId, srcBuffer, index, length) -> {
                byte[] message = new byte[length];
                srcBuffer.getBytes(index, message);
                System.out.println("Consumer <- " + new String(message));
                countDownLatch.countDown();
            });
        }
        long end = System.currentTimeMillis();
        System.out.println("cost time " + (end - start));
    }
}
public class App {
    private static final OneToOneRingBuffer BUFFER = new OneToOneRingBuffer(new UnsafeBuffer(
            ByteBuffer.allocate(1024 + RingBufferDescriptor.TRAILER_LENGTH)));

    public static void main(String[] args) throws Exception {
        ExecutorService executor = Executors.newCachedThreadPool();
        executor.execute(new Producer(BUFFER));
        executor.execute(new Consumer(BUFFER));
        Thread.currentThread().join();
    }
}

问题解答

1. OneToOneRingBuffer的线程安全性

OneToOneRingBuffer是线程安全的,它是Agrona专门为单生产者单消费者场景设计的无锁环形缓冲区,底层通过内存屏障保证跨线程的内存可见性,不需要额外同步锁就能安全地在单生产者和单消费者之间传递消息。

2. 生产消费异常的根本原因

你的问题和线程安全无关,核心是环形缓冲区容量不足:

  • 你创建的缓冲区仅分配了1024 + RingBufferDescriptor.TRAILER_LENGTH字节,其中1024字节是可用消息存储空间,剩余部分是Agrona内部管理用的trailer。
  • 当Config.SIZE=10_000时,生产者写入速度和消费者消费速度刚好匹配,缓冲区不会被填满,程序正常运行;但当SIZE=100_000时,生产者快速将缓冲区写满,后续buffer.write()调用会阻塞等待缓冲区腾出空间(OneToOneRingBuffer的write()默认是阻塞式),消费者来不及消费足够消息释放空间,最终导致生产者卡住,消费者无法读到新消息,CountDownLatch永远无法归0。

3. 修复方案

  • 扩容缓冲区:根据单条消息平均大小 × 预期最大待处理消息数,计算所需总容量,再加上RingBufferDescriptor.TRAILER_LENGTH。比如每条消息平均20字节,100_000条消息需要的容量至少是20 * 100000 + RingBufferDescriptor.TRAILER_LENGTH字节。
  • 处理写阻塞/失败:可以使用tryWrite()替代write(),该方法尝试写入失败时返回false,你可以在代码中处理这种情况,比如短暂等待后重试,避免生产者无限阻塞。

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

相关产品推荐
方舟 Agent Plan

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

最近更新时间:2026.08.24 19:33:23