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

Spring Kafka Streams迁K8s后RocksDB Manifest文件缺失故障求助

问题:Kafka Streams + RocksDB 在EFS存储下的状态文件异常与自动恢复需求

背景

我们运行一个基于Spring Kafka Stream Binder的有状态应用,内置RocksDB作为本地状态存储,从Kafka Topic同步状态数据。该应用在PCF环境使用本地磁盘存储时运行稳定,但迁移到AWS EKS后采用StatefulSet部署,使用AWS EFS通过StorageClass/PV/PVC挂载作为状态存储目录,应用运行2天到2周后会突然停止工作。

核心异常

应用崩溃的核心原因是RocksDB的Manifest文件缺失,触发的异常如下:

Encountered the following exception during processing and the registered exception handler opted to SHUTDOWN_CLIENT. The streams client is going to shut down now.  org.rocksdb.RocksDBException: While opening a file for sequentially reading: /mnt/data/appid/0_0/rocksdb/topic-STATE-STORE-0000000001/MANIFEST-000015: No such file or directory
            at org.rocksdb.RocksDB.open(Native Method) ~[rocksdbjni-6.19.3.jar:]
            at org.rocksdb.RocksDB.open(RocksDB.java:306) ~[rocksdbjni-6.19.3.jar:]
            at org.apache.kafka.streams.state.internals.RocksDBTimestampedStore.openRocksDB(RocksDBTimestampedStore.java:75) ~[kafka-streams-3.0.0.jar:]
            ... 18 more
        Wrapped by: org.apache.kafka.streams.errors.ProcessorStateException: Error opening store topic-STATE-STORE-0000000001 at location /mnt/data/appid/0_0/rocksdb/topic-STATE-STORE-0000000001
            at org.apache.kafka.streams.state.internals.RocksDBTimestampedStore.openRocksDB(RocksDBTimestampedStore.java:87) ~[kafka-streams-3.0.0.jar:]
            at org.apache.kafka.streams.state.internals.RocksDBStore.openDB(RocksDBStore.java:183) ~[kafka-streams-3.0.0.jar:]
            at org.apache.kafka.streams.state.internals.RocksDBStore.init(RocksDBStore.java:250) ~[kafka-streams-3.0.0.jar:]
            at org.apache.kafka.streams.state.internals.WrappedStateStore.init(WrappedStateStore.java:55) ~[kafka-streams-3.0.0.jar:]
            at org.apache.kafka.streams.state.internals.ChangeLoggingKeyValueBytesStore.init(ChangeLoggingKeyValueBytesStore.java:56) ~[kafka-streams-3.0.0.jar:]
            at org.apache.kafka.streams.state.internals.WrappedStateStore.init(WrappedStateStore.java:55) ~[kafka-streams-3.0.0.jar:]
            at org.apache.kafka.streams.state.internals.CachingKeyValueStore.init(CachingKeyValueStore.java:75) ~[kafka-streams-3.0.0.jar:]
            at org.apache.kafka.streams.state.internals.WrappedStateStore.init(WrappedStateStore.java:55) ~[kafka-streams-3.0.0.jar:]
            at org.apache.kafka.streams.state.internals.MeteredKeyValueStore.lambdainit1(MeteredKeyValueStore.java:126) ~[kafka-streams-3.0.0.jar:]
            at org.apache.kafka.streams.processor.internals.metrics.StreamsMetricsImpl.maybeMeasureLatency(StreamsMetricsImpl.java:769) ~[kafka-streams-3.0.0.jar:]
            at org.apache.kafka.streams.state.internals.MeteredKeyValueStore.init(MeteredKeyValueStore.java:126) ~[kafka-streams-3.0.0.jar:]
            at org.apache.kafka.streams.processor.internals.ProcessorStateManager.registerStateStores(ProcessorStateManager.java:201) ~[kafka-streams-3.0.0.jar:]
            at org.apache.kafka.streams.processor.internals.StateManagerUtil.registerStateStores(StateManagerUtil.java:97) ~[kafka-streams-3.0.0.jar:]
            at org.apache.kafka.streams.processor.internals.StandbyTask.initializeIfNeeded(StandbyTask.java:93) ~[kafka-streams-3.0.0.jar:]
            at org.apache.kafka.streams.processor.internals.TaskManager.tryToCompleteRestoration(TaskManager.java:436) ~[kafka-streams-3.0.0.jar:]
            at org.apache.kafka.streams.processor.internals.StreamThread.initializeAndRestorePhase(StreamThread.java:849) ~[kafka-streams-3.0.0.jar:]
            at org.apache.kafka.streams.processor.internals.StreamThread.runOnce(StreamThread.java:731) ~[kafka-streams-3.0.0.jar:]
            at org.apache.kafka.streams.processor.internals.StreamThread.runLoop(StreamThread.java:583) ~[kafka-streams-3.0.0.jar:]
            at org.apache.kafka.streams.processor.internals.StreamThread.run(StreamThread.java:555) [kafka-streams-3.0.0.jar:]

关联异常

同时出现的其他相关异常:

  • 文件重命名失败:
Encountered the following exception during processing and the registered exception handler opted to SHUTDOWN_CLIENT. The streams client is going to shut down now. org.rocksdb.RocksDBException: While renaming a file to /mnt/data/rto-examiner-v3/0_2/rocksdb/cf_rto_aggregated-STATE-STORE-0000000001/CURRENT: /mnt/data/rto-examiner-v3/0_2/rocksdb/cf_rto_aggregated-STATE-STORE-0000000001/000001.dbtmp: No such file or directory at org.rocksdb.RocksDB.open(Native Method) ~[rocksdbjni-6.19.3.jar:] at org.rocksdb.RocksDB.open(RocksDB.java:306)
  • 锁文件获取失败:
Encountered the following exception during processing and the registered exception handler opted to SHUTDOWN_CLIENT. The streams client is going to shut down now. org.rocksdb.RocksDBException: While lock file: /mnt/data/rto-examiner-v3/0_4/rocksdb/cf_rto_aggregated-STATE-STORE-0000000001/LOCK: Resource temporarily unavailable at org.rocksdb.RocksDB.open(Native Method) ~[rocksdbjni-6.19.3.jar:] at org.rocksdb.RocksDB.open(RocksDB.java:306) ~[rocksdbjni-6.19.3.jar:] at org.apache.kafka.streams.state.internals.RocksDBTimestampedStore.openRocksDB

故障表现

所有异常集中在同一时间段爆发,Stream客户端关闭后,所有消费者陷入无限循环尝试加入消费者组但无法成功,偶尔6个实例中有一个能成功但很快重启。

技术栈

  • Spring Boot 2.6.x
  • Spring Kafka Streams 3.0
  • Kubernetes:AWS EKS
  • 存储:AWS EFS(通过StorageClass/PV/PVC挂载)
  • 状态处理器:基于Spring Functions Topology定义

当前临时方案与诉求

目前只能通过修改消费者组名解决,但需要手动管理偏移量,操作繁琐。希望能通过编程方式捕获RocksDBException,自动重建或清理损坏的状态存储,让应用自动恢复。

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

相关产品推荐
方舟 Agent Plan

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

最近更新时间:2026.06.21 22:55:58