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

Broker滚动升级时Kafka Streams应用进入ERROR状态问题排查

问题:AWS MSK滚动升级后Kafka Streams应用进入ERROR状态

环境与现象

  • 使用Kafka 2.8.1(AWS MSK集群)
  • AWS执行Broker滚动升级后,Kafka Streams应用直接进入ERROR状态,查询状态存储时持续抛出以下异常:
java.lang.IllegalStateException: Error when retrieving state store.
    at org.springframework.cloud.stream.binder.kafka.streams.InteractiveQueryService.lambda$getHostInfo$3(InteractiveQueryService.java:237)
    at org.springframework.retry.support.RetryTemplate.doExecute(RetryTemplate.java:329)
    at org.springframework.retry.support.RetryTemplate.execute(RetryTemplate.java:209)
    at org.springframework.cloud.stream.binder.kafka.streams.InteractiveQueryService.getHostInfo(InteractiveQueryService.java:222)
    at com.sp.gos.processors.GossiperKVStoreQueryService.getHostInfo(GossiperKVStoreQueryService.java:108)
    at com.sp.gos.processors.GossiperKVStoreQueryService.get(GossiperKVStoreQueryService.java:68)
...
Caused by: java.lang.IllegalStateException: KafkaStreams is not running. State is ERROR.
    at org.apache.kafka.streams.KafkaStreams.validateIsRunningOrRebalancing(KafkaStreams.java:381)
    at org.apache.kafka.streams.KafkaStreams.queryMetadataForKey(KafkaStreams.java:1663)
    at org.springframework.cloud.stream.binder.kafka.streams.InteractiveQueryService.lambda$getHostInfo$2(InteractiveQueryService.java:227)
    ...

已配置的异常处理器(未触发)

已添加日志并继续的异常处理器配置,但遇到问题时未生效:

props.put(
        "default.deserialization.exception.handler",
        LogAndContinueExceptionHandler.class);
props.put(
        "default.production.exception.handler",
        LogAndContinueProductionExceptionHandler.class);

依赖版本

  • org.springframework.kafka:spring-kafka:jar - 3.0.9
  • org.apache.kafka:kafka-clients - 3.5.1
  • org.springframework.cloud:spring-cloud-stream-binder-kafka-streams - 4.0.4
  • org.apache.kafka:kafka-streams - 3.5.1

问题分析

1. 异常处理器未触发的原因

你配置的default.deserialization.exception.handler和default.production.exception.handler仅覆盖消息反序列化和生产消息阶段的异常,而当前问题是Kafka Streams实例本身因集群变更(Broker升级)进入ERROR状态,属于集群连接、元数据同步类的顶层异常,不在这两个处理器的处理范围内。

2. 根本原因

  • 跨版本兼容性问题:你使用的Kafka Streams(3.5.1)和Broker(2.8.1)版本跨度较大,高版本客户端与低版本Broker交互时,在集群滚动升级的场景下可能出现未预期的元数据同步失败、连接超时等问题,直接导致流实例进入ERROR状态。
  • 默认状态转换策略:Kafka Streams默认遇到无法自动恢复的集群异常时,会终止流处理并进入ERROR状态,不会自动尝试重启。

推荐处理方式

1. 修复版本兼容性

将Kafka Streams和kafka-clients版本降级到与Broker(2.8.1)匹配的2.8.x系列,比如2.8.1或2.8.2,彻底解决跨版本交互的兼容性问题。

2. 优化Kafka Streams配置

添加以下配置,增强集群变更时的容错能力:

# 增加连接重试的退避时间,给Broker升级恢复留出时间
retry.backoff.ms=1000
retry.backoff.max.ms=30000
# 配置状态存储的清理延迟,避免状态文件被误清理
kafka.streams.state.cleanup.delay.ms=86400000
# 确保交互式查询能正确定位实例
application.server=your-app-host:port

3. 添加状态监听器实现自动重启

注册StateListener监听流实例状态,当进入ERROR状态时自动尝试重启:

kafkaStreams.setStateListener((newState, oldState) -> {
    if (newState == KafkaStreams.State.ERROR) {
        new Thread(() -> {
            try {
                // 等待30秒,让Broker集群稳定后再重启
                Thread.sleep(30000);
                if (kafkaStreams.state() != KafkaStreams.State.RUNNING) {
                    kafkaStreams.close();
                    kafkaStreams.start();
                }
            } catch (InterruptedException e) {
                Thread.currentThread().interrupt();
            }
        }).start();
    }
});

4. 增强交互式查询的容错逻辑

在查询状态存储前,先检查Kafka Streams实例状态,避免抛出未处理的异常:

public Object get(String key) {
    KafkaStreams.State currentState = kafkaStreams.state();
    if (currentState != KafkaStreams.State.RUNNING && currentState != KafkaStreams.State.REBALANCING) {
        throw new IllegalStateException("流服务未就绪,请稍后重试");
    }
    // 执行状态存储查询逻辑
    // ...
}

5. MSK升级前的运维优化

升级Broker前,可分批暂停Kafka Streams应用,待Broker集群升级完成并稳定后,再逐步启动应用,减少集群变更对业务的冲击。


内容的提问来源于stack exchange,提问作者tuk

相关产品推荐
方舟 Agent Plan

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

最近更新时间:2026.07.10 17:43:20