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

