基于LinkedBlockingQueue实现的自定义Kafka Consumer的size()方法返回值不符合预期的问题求助
基于LinkedBlockingQueue实现的自定义Kafka Consumer的size()方法返回值不符合预期的问题求助
我在测试场景中发现原生的KafkaConsumer用起来不太顺手,于是基于LinkedBlockingQueue自己实现了一个封装类,技术栈用的是Spring Boot、Java 21、Guava和7.8.2版本的Kafka Testcontainers。这个自定义Consumer的核心功能(比如take()和poll()获取消息)都正常,但唯独在测试时先断言size()等于1的时候总是失败,明明之后能正常拿到消息,这让我有点困惑。
自定义KafkaConsumerCashOut实现
@Slf4j @Component public class KafkaConsumerCashOut<K, V> { private final BlockingQueue<KeyValue<K, V>> queue = new LinkedBlockingQueue<>(); private final KafkaConsumer<K, V> consumer; private final String topic; private final AtomicBoolean isShutdown = new AtomicBoolean(false); private final ExecutorService poller = Executors.newSingleThreadExecutor(); public KafkaConsumerCashOut(TestKafkaConsumerProperties kafkaStreamProperties, CashOutTopicProperties cashOutTopicProperties) { this.topic = cashOutTopicProperties.getTopic(); this.consumer = new KafkaConsumer<>(new HashMap<>(kafkaStreamProperties.getProperties())); } @SneakyThrows public KeyValue<K, V> take() { return queue.take(); } @SneakyThrows public @Nullable KeyValue<K, V> poll(Duration timeout) { return queue.poll(timeout.toMillis(), TimeUnit.MILLISECONDS); } public int size() { return queue.size(); } @PreDestroy public void destroy() throws InterruptedException { log.info("Closing consumer ..."); isShutdown.compareAndSet(false, true); poller.awaitTermination(5, TimeUnit.SECONDS); } @PostConstruct void init() { poller.submit(() -> { try { consumer.subscribe(Collections.singleton(topic)); while (!isShutdown.get()) { var polled = consumer.poll(Duration.ofMillis(100)); if (!polled.isEmpty()) { log.info("Kafka consumer is polled records {}", polled.count()); for (var record : polled) { Preconditions.checkState(queue.offer(KeyValue.of(record.key(), record.value()))); } } } } catch (Throwable th) { log.warn("Error in consumer", th); } finally { log.info("Consumer loop terminated"); try { consumer.close(); } catch (Throwable th) { log.warn("Error closing consumer", th); } } }); } }
KeyValue封装类
@Data @RequiredArgsConstructor(staticName = "of") public class KeyValue<K, V> { private final K key; private final V message; }
测试代码片段
@Autowired KafkaConsumerCashOut<String, WatchResult> cashOut; // ... 发送消息到Kafka的逻辑 ... // 这里断言总是失败 assertThat(cashOut.size()).isEqualTo(1); // 但这行能正常拿到消息,断言通过 var result = cashOut.take(); assertThat(result.getKey()).isEqualTo(OUTCOME_ID_STRING); assertThat(result.getMessage()).isEqualTo(expectedResult);
问题分析
核心原因是线程时序不匹配:
- 测试主线程在发送消息后立刻调用
size(),但此时后台的poller线程可能还没完成Kafka的poll()操作,或者还没把消息从poll结果中写入到LinkedBlockingQueue里。 - LinkedBlockingQueue的
size()是即时快照,当主线程执行断言时,队列里可能还没有被写入消息,所以返回0导致断言失败;而后续调用take()时,主线程会阻塞直到队列有元素,这时候后台线程已经完成了消息写入,所以能正常拿到数据。
解决方案
方案1:带超时的轮询断言(测试场景推荐)
在测试中不要直接断言size(),而是用轮询的方式等待队列达到预期大小,兼容Kafka消息传递和后台线程处理的延迟:
// 轮询等待队列size变为1,最多等5秒 assertThat(pollUntil(() -> cashOut.size() == 1, Duration.ofSeconds(5))) .withFailMessage("消息未在超时时间内进入队列") .isTrue(); // 通用轮询工具方法 private boolean pollUntil(BooleanSupplier condition, Duration timeout) { long endTime = System.currentTimeMillis() + timeout.toMillis(); while (System.currentTimeMillis() < endTime) { if (condition.getAsBoolean()) { return true; } try { Thread.sleep(100); } catch (InterruptedException e) { Thread.currentThread().interrupt(); return false; } } return false; }
方案2:给自定义Consumer添加等待size的方法
在KafkaConsumerCashOut中新增方法,让调用者可以阻塞等待队列达到指定大小:
@SneakyThrows public void awaitSize(int expectedSize, Duration timeout) { long endTime = System.currentTimeMillis() + timeout.toMillis(); while (queue.size() < expectedSize && System.currentTimeMillis() < endTime) { Thread.sleep(50); } if (queue.size() < expectedSize) { throw new TimeoutException("队列大小未在超时时间内达到预期值: " + expectedSize); } }
测试时调用:
cashOut.awaitSize(1, Duration.ofSeconds(5)); assertThat(cashOut.size()).isEqualTo(1);
内容来源于stack exchange
相关产品推荐
相关产品推荐

