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

