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

Spring Boot优雅关闭时Kafka状态存储获取失败问题咨询

背景与分析

我们的Spring Boot应用依赖Kafka存储的信息响应REST请求,通过InteractiveQueryService从全局状态存储获取数据。

但在优雅关闭过程中,正在处理的请求会触发如下错误:

java.lang.IllegalStateException: Error retrieving state store: my-global-store-name-v0
    at org.springframework.cloud.stream.binder.kafka.streams.InteractiveQueryService.lambda$getQueryableStore$1(InteractiveQueryService.java:153)
    at org.springframework.retry.support.RetryTemplate.doExecute(RetryTemplate.java:344)
    at org.springframework.retry.support.RetryTemplate.execute(RetryTemplate.java:217)
    at org.springframework.cloud.stream.binder.kafka.streams.InteractiveQueryService.getQueryableStore(InteractiveQueryService.java:103)
    ...

分析发现,错误根源在于:负责优雅关闭的WebServerGracefulShutdownLifecycle的phase值为Integer.MAX_VALUE - 1024,而StreamsBuilderFactoryManager的phase值为Integer.MAX_VALUE - 100。由于Spring生命周期中phase值越大,关闭顺序越早,导致StreamsBuilderFactoryManager在Web服务器开始优雅关闭前就已关闭,此时请求无法访问全局状态存储。

日志可验证该关闭顺序:

2025-01-08 17:04:16.225 DEBUG  org.springframework.context.support.DefaultLifecycleProcessor [Stopping beans in phase 2147483547]
...
2025-01-08 17:04:16.572 DEBUG  org.apache.kafka.streams.processor.internals.GlobalStateManagerImpl [Closing global storage engine my-global-store-name-v0]
...
2025-01-08 17:04:16.577 DEBUG  org.springframework.context.support.DefaultLifecycleProcessor [Bean 'streamsBuilderFactoryManager' completed its stop procedure]
...
2025-01-08 17:04:16.584 DEBUG  org.springframework.context.support.DefaultLifecycleProcessor [Stopping beans in phase 2147482623]  
2025-01-08 17:04:16.584 INFO   org.springframework.boot.web.embedded.tomcat.GracefulShutdown [Commencing graceful shutdown. Waiting for active requests to complete]  
2025-01-08 17:04:16.587 INFO   org.springframework.boot.web.embedded.tomcat.GracefulShutdown [Graceful shutdown complete]  

StreamsBuilderFactoryManager的Javadoc显示,选择接近Integer.MAX_VALUE的phase值是有意为之:

This {@link SmartLifecycle} class ensures that the bean created from it is started very late through the bootstrap process by setting the phase value closer to Integer.MAX_VALUE. This is to guarantee that the {@link StreamsBuilderFactoryBean} on a function with multiple bindings is only started after all the binding phases have completed successfully.

此外,值同为Integer.MAX_VALUE - 100的常量AbstractMessageListenerContainer.DEFAULT_PHASE有如下说明(暂未验证):

// The default org.springframework.context.SmartLifecycle phase for listener containers 2147483547.

问题

基于上述分析,StreamsBuilderFactoryManager的phase值高于WebServerGracefulShutdownLifecycle是否属于Bug?若是,我将在GitHub提交Issue并修复;若不是,如何确保优雅关闭时请求仍能访问全局状态存储?


解决方案

1. 是否属于Bug

这属于设计缺陷:StreamsBuilderFactoryManager设置高phase值是为了保证启动顺序,但未考虑Spring生命周期的反向逻辑——启动phase越高,关闭顺序越早。对于Web请求依赖Kafka Streams状态存储的场景,当前配置会导致优雅关闭期间服务不可用,因此有必要提交Issue请求调整默认phase值,使其低于WebServerGracefulShutdownLifecycle的phase(即Integer.MAX_VALUE - 1024),确保Web服务器完成优雅关闭后再关闭Kafka Streams相关组件。

2. 临时修复方案(无需修改框架源码)

方案A:调整StreamsBuilderFactoryManager的phase值

通过Bean后置处理器修改目标Bean的phase值,确保其关闭顺序晚于Web服务器优雅关闭:

import org.springframework.beans.BeansException;
import org.springframework.beans.factory.config.BeanPostProcessor;
import org.springframework.cloud.stream.binder.kafka.streams.StreamsBuilderFactoryManager;
import org.springframework.context.SmartLifecycle;
import org.springframework.stereotype.Component;

@Component
public class StreamsBuilderFactoryManagerPhaseAdjuster implements BeanPostProcessor {

    private static final int WEB_SERVER_SHUTDOWN_PHASE = Integer.MAX_VALUE - 1024;

    @Override
    public Object postProcessAfterInitialization(Object bean, String beanName) throws BeansException {
        if (bean instanceof StreamsBuilderFactoryManager) {
            ((StreamsBuilderFactoryManager) bean).setPhase(WEB_SERVER_SHUTDOWN_PHASE - 1);
        }
        return bean;
    }
}
方案B:自定义生命周期控制Bean

创建自定义SmartLifecycleBean,依赖StreamsBuilderFactoryManager并设置更低的phase值,确保Web服务器关闭后再触发Kafka Streams组件的关闭:

import org.springframework.cloud.stream.binder.kafka.streams.StreamsBuilderFactoryManager;
import org.springframework.context.SmartLifecycle;
import org.springframework.stereotype.Component;

@Component
public class KafkaStreamsGracefulShutdownLifecycle implements SmartLifecycle {

    private final StreamsBuilderFactoryManager streamsManager;
    private boolean running = false;

    public KafkaStreamsGracefulShutdownLifecycle(StreamsBuilderFactoryManager streamsManager) {
        this.streamsManager = streamsManager;
    }

    @Override
    public void start() {
        running = true;
    }

    @Override
    public void stop() {
        streamsManager.stop();
    }

    @Override
    public boolean isRunning() {
        return running;
    }

    @Override
    public int getPhase() {
        // 比WebServerGracefulShutdownLifecycle的phase低,确保最后关闭
        return Integer.MAX_VALUE - 1025;
    }
}

同时需禁用原StreamsBuilderFactoryManager的自动关闭:

spring.cloud.stream.kafka.streams.binder.auto-startup=false

3. 提交Issue建议

在Spring Cloud Stream或Spring Boot的GitHub仓库提交Issue时,需包含:

  • 场景描述:Web请求依赖Kafka Streams全局状态存储,优雅关闭期间触发状态存储不可访问错误
  • 问题根源:phase值导致的关闭顺序倒置
  • 日志片段:展示关闭顺序的关键日志
  • 修复建议:调整StreamsBuilderFactoryManager的默认phase值,使其低于WebServerGracefulShutdownLifecycle的phase

内容的提问来源于stack exchange,提问作者Cédric Schaller

相关产品推荐
方舟 Agent Plan

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

最近更新时间:2026.06.15 01:53:13