如何在测试用例中等待所有Disruptor消息被消费?
Disruptor单元测试中等待消费完成的简便方法
针对Disruptor单元测试里等待所有EventHandler完成消息消费的需求,有几种简便实用的实现方式:
1. 使用CountDownLatch(最直接常用)
初始化一个计数等于待发布消息数量的CountDownLatch,在每个EventHandler的onEvent方法处理完消息后调用countDown(),测试主线程通过latch.await()等待所有消费完成,建议加上超时时间避免测试无限挂起。
示例代码:
// 测试类中定义 int messageCount = 10; CountDownLatch consumeLatch = new CountDownLatch(messageCount); // 配置EventHandler EventHandler<MyEvent> eventHandler = (event, sequence, endOfBatch) -> { // 执行业务消费逻辑 consumeLatch.countDown(); }; // 初始化Disruptor并绑定handler Disruptor<MyEvent> disruptor = new Disruptor<>(MyEvent::new, 1024, Executors.defaultThreadFactory()); disruptor.handleEventsWith(eventHandler); disruptor.start(); // 发布消息 RingBuffer<MyEvent> ringBuffer = disruptor.getRingBuffer(); for (int i = 0; i < messageCount; i++) { long sequence = ringBuffer.next(); ringBuffer.get(sequence).setData("消息" + i); ringBuffer.publish(sequence); } // 等待消费完成,超时5秒 consumeLatch.await(5, TimeUnit.SECONDS); // 执行断言操作 // ...
2. 利用Disruptor的序列跟踪(无侵入式)
Disruptor的每个EventHandler都会维护自身的消费序列,通过对比环形缓冲区的当前游标(已发布消息的最后序列)与所有handler的消费序列,判断是否所有消息都被处理。这种方式不需要修改handler的业务代码,更干净。
示例代码:
// 封装等待工具方法 private void waitForAllConsumed(RingBuffer<?> ringBuffer, Sequence... handlerSequences) throws InterruptedException { long latestPublishedSeq = ringBuffer.getCursor(); // 循环等待所有handler的消费序列追上已发布的最后序列 for (Sequence seq : handlerSequences) { while (seq.get() < latestPublishedSeq) { Thread.sleep(10); // 短休眠避免空轮询占用CPU } } } // 测试中使用 Disruptor<MyEvent> disruptor = // 初始化配置... disruptor.start(); RingBuffer<MyEvent> ringBuffer = disruptor.getRingBuffer(); // 发布消息... // 获取所有EventHandler的消费序列 Sequence[] handlerSequences = disruptor.getConsumerRepository().getAllSequences(); // 等待消费完成 waitForAllConsumed(ringBuffer, handlerSequences); // 执行断言操作 // ...
3. 结合测试断言库(优雅的测试风格)
如果使用AssertJ这类断言库,可以利用其异步断言能力,直接断言消费后的结果状态(比如已消费消息列表的大小),框架会自动等待直到条件满足或超时。
示例代码:
// 定义线程安全的消费结果容器 List<String> consumedDataList = Collections.synchronizedList(new ArrayList<>()); // 配置EventHandler EventHandler<MyEvent> eventHandler = (event, sequence, endOfBatch) -> { consumedDataList.add(event.getData()); }; // 初始化Disruptor并发布消息... int messageCount = 10; // 发布消息逻辑... // 等待消费完成并直接断言,超时5秒 assertThat(consumedDataList) .await() .atMost(5, TimeUnit.SECONDS) .hasSize(messageCount);
注意事项
- 所有等待逻辑都建议加上超时时间,防止因消费逻辑异常导致测试无限挂起。
- 如果是多EventHandler并行消费的场景,要确保获取所有handler的序列进行判断,不要遗漏。
内容的提问来源于stack exchange,提问作者lance-java
相关产品推荐
相关产品推荐

