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
相关产品推荐
相关产品推荐

