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消费者组初始偏移量设置、主题动态创建后的元数据同步延迟相关:
- 销毁重建容器后,Kafka内部元数据(消费者组偏移量、主题分区信息)被清空,Quarkus的SmallRye Kafka客户端首次连接时,可能在元数据未完全同步的情况下发送消息,导致Broker接收消息但无法路由到未完成初始化的消费者。
- 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
验证步骤
- 销毁并重建Kafka容器:
docker-compose down && docker-compose up -d - 运行JUnit测试,检查首批消息是否正常接收
- 重复执行测试,确认问题不再出现
内容的提问来源于stack exchange,提问作者Slush
相关产品推荐
相关产品推荐

