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

Quarkus应用中Kafka、Redis、gRPC及CDI Bean实际启停顺序排查与关闭错误解决咨询

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的生命周期规则,启停顺序严格遵循「启动顺序的逆序」,具体到你的组件:

启动顺序(从早到晚):

  1. 核心CDI容器初始化
  2. gRPC客户端、Redis客户端初始化(因为CDI Bean会依赖它们,所以提前启动)
  3. Kafka消费者初始化(SmallRye Kafka扩展启动)
  4. 被依赖的CDI Bean启动(比如FeatureService、DataProvider)
  5. 依赖其他Bean的CDI Bean启动(比如BusinessService,依赖前面的几个Bean)
  6. 标记@Startup(n)的Bean启动(比如你的CachePreloadService,n值越小启动越早)

关闭顺序(从早到晚,启动顺序的完全逆序):

  1. 标记@Startup(n)的Bean先销毁(CachePreloadService)
  2. 依赖其他Bean的CDI Bean销毁(BusinessService)
  3. 被依赖的CDI Bean销毁(FeatureService、DataProvider)
  4. Kafka消费者停止(这里要注意:如果有正在处理的异步消息,默认优雅关闭会等处理完成,但如果你的Uni没有正确绑定生命周期,可能会出现延迟)
  5. Redis客户端关闭
  6. gRPC客户端关闭
  7. 核心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

相关产品推荐
方舟 Agent Plan

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

最近更新时间:2026.04.07 09:58:08