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

基于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);

问题分析

核心原因是线程时序不匹配:

  1. 测试主线程在发送消息后立刻调用size(),但此时后台的poller线程可能还没完成Kafka的poll()操作,或者还没把消息从poll结果中写入到LinkedBlockingQueue里。
  2. 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

相关产品推荐
方舟 Agent Plan

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

最近更新时间:2026.04.07 11:19:36