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

Quarkus重启Docker容器后丢失首批SmallRye Kafka消息求助

Quarkus Kafka消费者在重建Kafka容器后丢失首批测试消息

问题描述

使用Confluentinc Kafka(ZooKeeper/Broker)镜像部署服务,每次销毁并重建Docker容器后,Quarkus应用的JUnit测试发送的首批消息会丢失,重复执行测试则无此问题。疑似问题出在Kafka Broker与Quarkus注入的@ApplicationScoped监听器Stuff之间。测试前用Offset Explorer连接集群有时能避免该故障。

已完成的排查操作

  • 通过@Observes StartupEvent强制预加载Stuff监听器,仍无法接收消息
  • 非Quarkus终端发送的消息可被纯Kafka消费者正常接收
  • @Incoming监听方法的日志无输出,说明消息未到达消费者
  • 一次性发送1000条消息时,首批消息既未被投递也不在Kafka队列(Kafka显示已确认接收)
  • 容器创建后等待3分钟再执行测试,问题依旧
  • 多主题发送时仅首批消息丢失,第二个主题可正常接收
  • 最小复现示例中添加Schema Registry可临时解决,但正式应用无效
  • 仅重启容器不会触发问题,必须销毁重建才会复现

核心原因分析

该问题主要和Kafka消费者组初始偏移量设置、主题动态创建后的元数据同步延迟相关:

  1. 销毁重建容器后,Kafka内部元数据(消费者组偏移量、主题分区信息)被清空,Quarkus的SmallRye Kafka客户端首次连接时,可能在元数据未完全同步的情况下发送消息,导致Broker接收消息但无法路由到未完成初始化的消费者。
  2. Offset Explorer连接时会主动触发Kafka元数据同步,提前完成消费者组和主题的初始化,因此后续测试正常。

解决方案

1. 配置消费者初始偏移量

在application.properties中为消费者添加偏移量配置,确保首次消费从最开始读取:

mp.messaging.incoming.TestChannel1.auto.offset.reset=earliest
mp.messaging.incoming.TestChannel2.auto.offset.reset=earliest

默认auto.offset.reset为latest,当消费者组无历史偏移量时会从最新消息开始消费,首批消息可能在消费者初始化完成前发送,导致丢失。

2. 提前创建测试主题

在Docker Compose的Kafka Broker配置中添加自动创建主题的参数,确保测试主题在容器启动时就存在,避免动态创建的元数据延迟:

environment:
  # 保留原有配置
  KAFKA_AUTO_CREATE_TOPICS_ENABLE: true
  KAFKA_CREATE_TOPICS: "TestChannel1:1:1,TestChannel2:1:1"

格式为主题名:分区数:副本数,单Broker环境下副本数设为1。

3. 测试前等待消费者初始化完成

在Quarkus测试中添加等待逻辑,确保消费者完全连接并订阅主题后再发送消息,可利用SmallRye的ConnectorStatus检查:

@Inject
ConnectorStatus connectorStatus;

@BeforeEach
void waitForConsumerReady() throws InterruptedException {
    // 等待TestChannel1消费者就绪
    while (!connectorStatus.isReady("TestChannel1")) {
        Thread.sleep(100);
    }
}

需要添加依赖:

<dependency>
    <groupId>io.smallrye.reactive</groupId>
    <artifactId>smallrye-reactive-messaging-health</artifactId>
</dependency>

4. 调整消费者组初始化延迟

在Docker Compose的Kafka Broker配置中,调大KAFKA_GROUP_INITIAL_REBALANCE_DELAY_MS,给消费者组足够初始化时间:

environment:
  # 保留原有配置
  KAFKA_GROUP_INITIAL_REBALANCE_DELAY_MS: 1000

5. 强化生产者消息确认级别

调整生产者acks配置,确保消息被Broker持久化后才确认:

mp.messaging.outgoing.TestChannel1-out.acks=all
mp.messaging.outgoing.TestChannel2-out.acks=all

验证步骤

  1. 销毁并重建Kafka容器:docker-compose down && docker-compose up -d
  2. 运行JUnit测试,检查首批消息是否正常接收
  3. 重复执行测试,确认问题不再出现

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

相关产品推荐
方舟 Agent Plan

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

最近更新时间:2026.08.08 09:40:45