使用RocksDB将Flink键控状态写入Azure UltraSSD LRS时触发异常
问题描述
- 在Kubernetes Pod中挂载Azure UltraSSD_LRS类型PVC,部署Flink应用存储键控状态
- Flink作业配置:
state.backend: rocksdb state.backend.rocksdb.localdir: "/data/flink/state/local state.checkpoints.dir: abfs://XXXXXXX/flink-checkpointing/
- 代码中设置DB存储路径:
EmbeddedRocksDBStateBackend embeddedRocksDb = new EmbeddedRocksDBStateBackend(true); embeddedRocksDb.setDbStoragePath("file:///data/flink/state/job/"); env.setStateBackend(embeddedRocksDb); env.enableCheckpointing(5000);
- 现象:磁盘挂载正常且可写入数据,但运行一段时间后抛出RocksDB异常
错误日志
Caused by: org.apache.flink.util.SerializedThrowable: org.rocksdb.RocksDBException: while link file to /data/flink/state/job/job_cb48da41cb620a68170593dc09789f2b_op_KeyedCoProcessOperator_1c449adacf198fac8b664046293f1fdf__2_4__uuid_4c6051b6-b588-4bcb-a740-fe4ddb01beae/chk-8.tmp/000014.sst: /data/flink/state/job/job_cb48da41cb620a68170593dc09789f2b_op_KeyedCoProcessOperator_1c449adacf198fac8b664046293f1fdf__2_4__uuid_4c6051b6-b588-4bcb-a740-fe4ddb01beae/db/000014.sst: Operation not supported at org.rocksdb.Checkpoint.createCheckpoint(Native Method) ~[flink-dist-1.17.1.jar:1.17.1] at org.rocksdb.Checkpoint.createCheckpoint(Checkpoint.java:51) ~[flink-dist-1.17.1.jar:1.17.1] at org.apache.flink.contrib.streaming.state.snapshot.RocksDBSnapshotStrategyBase.takeDBNativeCheckpoint(RocksDBSnapshotStrategyBase.java:170) ~[flink-dist-1.17.1.jar:1.17.1] at org.apache.flink.contrib.streaming.state.snapshot.RocksDBSnapshotStrategyBase.syncPrepareResources(RocksDBSnapshotStrategyBase.java:156) ~[flink-dist-1.17.1.jar:1.17.1] at org.apache.flink.contrib.streaming.state.snapshot.RocksDBSnapshotStrategyBase.syncPrepareResources(RocksDBSnapshotStrategyBase.java:76) ~[flink-dist-1.17.1.jar:1.17.1] at org.apache.flink.runtime.state.SnapshotStrategyRunner.snapshot(SnapshotStrategyRunner.java:77) ~[flink-dist-1.17.1.jar:1.17.1] at org.apache.flink.contrib.streaming.state.RocksDBKeyedStateBackend.snapshot(RocksDBKeyedStateBackend.java:593) ~[flink-dist-1.17.1.jar:1.17.1] at org.apache.flink.streaming.api.operators.StreamOperatorStateHandler.snapshotState(StreamOperatorStateHandler.java:246) ~[flink-dist-1.17.1.jar:1.17.1] at org.apache.flink.streaming.api.operators.StreamOperatorStateHandler.snapshotState(StreamOperatorStateHandler.java:173) ~[flink-dist-1.17.1.jar:1.17.1] at org.apache.flink.streaming.api.operators.AbstractStreamOperator.snapshotState(AbstractStreamOperator.java:336) ~[flink-dist-1.17.1.jar:1.17.1] at org.apache.flink.streaming.runtime.tasks.RegularOperatorChain.checkpointStreamOperator(RegularOperatorChain.java:228) ~[flink-dist-1.17.1.jar:1.17.1] at org.apache.flink.streaming.runtime.tasks.RegularOperatorChain.buildOperatorSnapshotFutures(RegularOperatorChain.java:213) ~[flink-dist-1.17.1.jar:1.17.1] at org.apache.flink.streaming.runtime.tasks.RegularOperatorChain.snapshotState(RegularOperatorChain.java:192) ~[flink-dist-1.17.1.jar:1.17.1] at org.apache.flink.streaming.runtime.tasks.SubtaskCheckpointCoordinatorImpl.takeSnapshotSync(SubtaskCheckpointCoordinatorImpl.java:715) ~[flink-dist-1.17.1.jar:1.17.1] at org.apache.flink.streaming.runtime.tasks.SubtaskCheckpointCoordinatorImpl.checkpointState(SubtaskCheckpointCoordinatorImpl.java:350) ~[flink-dist-1.17.1.jar:1.17.1] at org.apache.flink.streaming.runtime.tasks.StreamTask.lambda$performCheckpoint$14(StreamTask.java:1299) ~[flink-dist-1.17.1.jar:1.17.1] at org.apache.flink.streaming.runtime.tasks.StreamTaskActionExecutor$1.runThrowing(StreamTaskActionExecutor.java:50) ~[flink-dist-1.17.1.jar:1.17.1] at org.apache.flink.streaming.runtime.tasks.StreamTask.performCheckpoint(StreamTask.java:1287) ~[flink-dist-1.17.1.jar:1.17.1] at org.apache.flink.streaming.runtime.tasks.StreamTask.triggerCheckpointOnBarrier(StreamTask.java:1244) ~[flink-dist-1.17.1.jar:1.17.1] ... 23 more
解决方案
这个错误的核心原因是RocksDB创建Checkpoint时默认使用硬链接来复制文件,但Azure UltraSSD_LRS的文件系统不支持硬链接操作。
有两种可行的解决方式:
- 禁用RocksDB的硬链接功能
在Flink配置文件中添加以下参数,强制RocksDB在生成Checkpoint时复制文件而非使用硬链接:
state.backend.rocksdb.checkpoint.use-hardlinks: false # 可选:设置文件复制线程数,提升Checkpoint生成效率 state.backend.rocksdb.checkpoint.transfer.thread.num: 4
或者直接在代码中配置:
embeddedRocksDb.setUseHardlinksForCheckpoints(false);
- 将RocksDB本地存储路径切换到Pod临时目录
把state.backend.rocksdb.localdir和代码中的dbStoragePath改为Pod的临时目录(如/tmp/flink/state),临时目录使用节点本地磁盘,支持硬链接操作。远程Checkpoint存储路径(abfs地址)保持不变,生成的Checkpoint依然会正常上传到远程存储。
注意:使用临时目录时,Pod重启后本地状态会丢失,但Flink会自动从远程Checkpoint恢复,不会影响作业的容错性。
内容的提问来源于stack exchange,提问作者Aman Vaishya
相关产品推荐
相关产品推荐

