Quarkus应用中Kafka、Redis、gRPC及CDI Bean实际启停顺序排查与关闭错误解决咨询
兄弟,我太懂你这种被Quarkus关闭顺序坑到的痛苦了!之前我维护的一个Quarkus项目也碰到过一模一样的问题:Kafka消费者还在跑,Redis已经被关了,一堆报错看得头大。结合你用的Quarkus 3.19.3的环境,给你几个实用的排查和解决方法:
一、怎么精准排查实际的关闭顺序?
别再盯着全量DEBUG日志找了,这几个方法直接帮你定位:
1. 给关键组件加生命周期埋点日志
最直接的方法就是在你关心的每个CDI Bean、依赖组件里加@PreDestroy钩子,打印明确的关闭日志:
import org.jboss.logging.Logger; import jakarta.annotation.PreDestroy; import jakarta.enterprise.context.ApplicationScoped; @ApplicationScoped public class FeatureService { private static final Logger log = Logger.getLogger(FeatureService.class); // 你的Redis相关代码... @PreDestroy void onShutdown() { log.info("=== 【SHUTDOWN】 {} 开始销毁 ===", this.getClass().getSimpleName()); } }
同样的,给MessageConsumer(Kafka)、DataProvider(gRPC)、BusinessService这些都加上,日志里就会清晰看到每个组件的关闭时间点,配合Quarkus默认的日志时间戳,就能直接排序出顺序。
2. 针对性开启组件生命周期日志
不用全开DEBUG,只开启和生命周期相关的日志分类,噪音瞬间减少:
启动时加这些JVM参数:
-Dquarkus.log.category."io.quarkus.arc".level=INFO -Dquarkus.log.category."io.quarkus.smallrye.kafka".level=INFO -Dquarkus.log.category."io.quarkus.redis".level=INFO -Dquarkus.log.category."io.quarkus.grpc".level=INFO
io.quarkus.arc是Quarkus CDI的核心包,会输出所有CDI Bean的销毁顺序- 后面几个分别对应Kafka、Redis、gRPC扩展的生命周期日志,能看到客户端的启动/关闭时间
3. 全局生命周期事件观察者
写一个全局的监听Bean,专门捕获各个组件的生命周期事件,不用给每个Bean都加钩子:
import org.jboss.logging.Logger; import jakarta.enterprise.context.ApplicationScoped; import jakarta.enterprise.event.Observes; import jakarta.interceptor.Interceptor; import io.quarkus.arc.runtime.BeanDestroyed; import io.quarkus.smallrye.kafka.runtime.KafkaConsumerStoppedEvent; import io.quarkus.redis.runtime.client.RedisClientStoppedEvent; import io.quarkus.grpc.runtime.GrpcServerStoppedEvent; @ApplicationScoped public class LifecycleTracker { private static final Logger log = Logger.getLogger(LifecycleTracker.class); // 监听所有CDI Bean的销毁事件 void trackBeanDestroy(@Observes @Priority(Interceptor.Priority.LIBRARY_BEFORE) BeanDestroyed<?> event) { log.info("【CDI Bean销毁】: {}", event.getBean().getBeanClass().getSimpleName()); } // 监听Kafka消费者关闭 void trackKafkaShutdown(@Observes KafkaConsumerStoppedEvent event) { log.info("【Kafka 关闭】: 消费者已停止"); } // 监听Redis客户端关闭 void trackRedisShutdown(@Observes RedisClientStoppedEvent event) { log.info("【Redis 关闭】: 客户端已停止"); } // 监听gRPC服务/客户端关闭 void trackGrpcShutdown(@Observes GrpcServerStoppedEvent event) { log.info("【gRPC 关闭】: 服务端已停止"); } }
这个Bean会自动捕获所有相关的关闭事件,日志里会按时间顺序输出各个组件的关闭节点,一目了然。
二、Quarkus中Kafka、Redis、gRPC、CDI Bean的默认启停顺序
根据Quarkus 3.x的生命周期规则,启停顺序严格遵循「启动顺序的逆序」,具体到你的组件:
启动顺序(从早到晚):
- 核心CDI容器初始化
- gRPC客户端、Redis客户端初始化(因为CDI Bean会依赖它们,所以提前启动)
- Kafka消费者初始化(SmallRye Kafka扩展启动)
- 被依赖的CDI Bean启动(比如
FeatureService、DataProvider) - 依赖其他Bean的CDI Bean启动(比如
BusinessService,依赖前面的几个Bean) - 标记
@Startup(n)的Bean启动(比如你的CachePreloadService,n值越小启动越早)
关闭顺序(从早到晚,启动顺序的完全逆序):
- 标记
@Startup(n)的Bean先销毁(CachePreloadService) - 依赖其他Bean的CDI Bean销毁(
BusinessService) - 被依赖的CDI Bean销毁(
FeatureService、DataProvider) - Kafka消费者停止(这里要注意:如果有正在处理的异步消息,默认优雅关闭会等处理完成,但如果你的
Uni没有正确绑定生命周期,可能会出现延迟) - Redis客户端关闭
- gRPC客户端关闭
- 核心CDI容器销毁
你碰到的「Kafka consumers processing after Redis is closed」错误,本质就是Kafka消费者的优雅关闭没有等所有消息处理完成,Redis客户端就已经被关闭了——因为Redis客户端的关闭顺序在Kafka消费者之后?不对,其实是你的消息处理逻辑用了FeatureService(依赖Redis),而FeatureService被销毁后Redis客户端才关闭,但如果Kafka的消息处理在FeatureService销毁后还在跑,就会报错。
三、怎么解决你的关闭错误?
根据上面的顺序,给你几个实操的修复方案:
1. 延长Kafka优雅关闭超时
让Quarkus等Kafka消费者处理完所有消息再继续后续的关闭流程,在application.properties里加:
quarkus.smallrye-kafka.graceful-shutdown-timeout=30s quarkus.smallrye-kafka.consumer.graceful-shutdown=true
这个参数会让Quarkus等待Kafka消费者完成所有正在处理的消息,再关闭消费者,避免后续依赖被销毁后还在处理消息。
2. 调整CDI Bean的销毁优先级
给依赖外部资源的Bean(比如FeatureService、DataProvider)设置更高的销毁优先级,让它们晚于Kafka消费者销毁:
@ApplicationScoped @Priority(Interceptor.Priority.APPLICATION_AFTER) // 晚于普通Bean销毁 public class FeatureService { // ... 你的代码 }
这样FeatureService会比其他普通CDI Bean晚销毁,保证Kafka消息处理时Redis还可用。
3. 手动控制Kafka消息处理的生命周期
在MessageConsumer里加逻辑,确保关闭时所有正在处理的消息都完成,再允许Bean销毁:
import org.jboss.logging.Logger; import jakarta.annotation.PreDestroy; import jakarta.enterprise.context.ApplicationScoped; import java.util.ArrayList; import java.util.Collections; import java.util.List; import java.util.concurrent.CompletableFuture; import java.util.concurrent.TimeUnit; import io.smallrye.mutiny.Uni; import io.vertx.core.json.JsonNode; @ApplicationScoped public class MessageConsumer { private static final Logger log = Logger.getLogger(MessageConsumer.class); private final List<Uni<Void>> activeProcesses = Collections.synchronizedList(new ArrayList<>()); private final CompletableFuture<Void> allProcessingComplete = new CompletableFuture<>(); @Incoming("input-topic") public Uni<Void> process(JsonNode msg) { Uni<Void> processingTask = doActualProcessing(msg) .onTermination(() -> activeProcesses.remove(processingTask)); activeProcesses.add(processingTask); return processingTask; } // 实际的消息处理逻辑 private Uni<Void> doActualProcessing(JsonNode msg) { // 你的业务代码,比如调用FeatureService操作Redis return Uni.createFrom().voidItem(); } @PreDestroy void onShutdown() { log.info("开始等待所有Kafka消息处理完成..."); // 等待所有正在处理的任务完成 Uni.combine().all().unis(activeProcesses) .discardItems() .subscribe().with( ignored -> allProcessingComplete.complete(null), error -> { log.error("消息处理出错", error); allProcessingComplete.completeExceptionally(error); } ); try { allProcessingComplete.get(30, TimeUnit.SECONDS); log.info("所有Kafka消息处理完成,准备关闭消费者"); } catch (Exception e) { log.error("等待消息处理超时,强制关闭", e); } } }
这个方法会在MessageConsumer销毁前,强制等待所有正在处理的Kafka消息完成,彻底避免依赖组件关闭后还在处理消息的情况。
内容来源于stack exchange

