news 2026/9/24 7:39:08

Flink Working Directory 完全指南:进程工作目录配置与跨重启本地恢复实战

作者头像

张小明

前端开发工程师

1.2k 24
文章封面图
Flink Working Directory 完全指南:进程工作目录配置与跨重启本地恢复实战
  • 大数据
  • 流处理
  • 批处理
  • 数据工程

【免费下载链接】flink

项目地址:https://gitcode.com/gh_mirrors/fli/flink
点击查看免费下载

Working Directory(工作目录)是 Flink 为 JobManager 与 TaskManager 进程提供的本地持久化目录,用于存放进程重启后可以恢复的运行时信息,其设计由 FLIP-198 引入、并由 FLIP-201 扩展出跨进程重启的本地恢复能力。本文以当前仓库(Apache Flink)Standalone 部署文档 working_directory.md 为主体,结合flink-runtime/flink-core的源码实现,系统讲解 Working Directory 的目录结构、配置项、进程资源 ID 的确定规则、存储在其中的三类制品,以及如何利用它实现 TaskManager 重启后从本地快速恢复状态。读完本文,你将能够为 Standalone 集群正确配置工作目录,并搭建一套"进程重启不丢失本地状态"的高可用实践方案。

什么是 Working Directory

Working Directory 是 Flink 进程(JobManager 和 TaskManager)用于存放进程重启后可以恢复的信息的本地目录。它的核心价值在于:当进程以相同身份(identity)重启、且仍然能够访问承载该目录的存储卷时,之前写入工作目录的数据可以继续被读取和使用,从而避免从远程存储重新拉取。

在源码中,工作目录由 WorkingDirectory.java 统一管理。该类的注释明确指出:

Class that manages a working directory for a process/instance. When being instantiated, this class makes sure that the specified working directory exists.

也就是说,WorkingDirectory在实例化时会确保目录存在(不存在则创建),并在其内部自动规划好一组固定的子目录结构。从构造逻辑看(WorkingDirectory.java),每个工作目录根下固定包含以下子目录:

子目录用途
tmp/进程临时文件目录,创建时会先清空(FileUtils.cleanDirectory
localState/本地状态目录,供本地恢复(local recovery)使用
blobStorage/Blob 存储目录,供 BlobServer / BlobCache 使用
slotAllocationSnapshots/槽位分配快照目录

这些子目录可以通过WorkingDirectory提供的getTmpDirectory()getLocalStateDirectory()getBlobStorageDirectory()getSlotAllocationSnapshotDirectory()等方法在运行时获取。

工作目录的目录结构

Flink 为两类进程分别规划了工作目录,其命名规则如下:

  • JobManager 工作目录:<WORKING_DIR_BASE>/jm_<JM_RESOURCE_ID>
  • TaskManager 工作目录:<WORKING_DIR_BASE>/tm_<TM_RESOURCE_ID>

其中:

  • <WORKING_DIR_BASE>是工作目录基路径(base);
  • <JM_RESOURCE_ID>是 JobManager 进程的资源 ID(resource id);
  • <TM_RESOURCE_ID>是 TaskManager 进程的资源 ID。

也就是说,同一个基路径下,不同进程通过jm_/tm_前缀与各自的资源 ID 区分目录。这一规则与源码中目录生成逻辑完全一致:在 ClusterEntrypointUtils.java 中,generateTaskManagerWorkingDirectoryFile使用"tm_" + resourceId作为目录名,generateJobManagerWorkingDirectoryFile使用"jm_" + resourceId作为目录名。

整个生成流程(generateWorkingDirectoryFile,见 ClusterEntrypointUtils.java)的决策顺序是:

  1. 若配置了进程专属的 working-dir 选项(如process.jobmanager.working-dir/process.taskmanager.working-dir),则直接以其为基路径;
  2. 否则若配置了通用选项process.working-dir,则以其为基路径(专属选项通过withFallbackKeys回退到通用选项);
  3. 若以上均未配置,则从io.tmp.dirs随机挑选一个临时目录作为基路径(对应ConfigurationUtils.getRandomTempDirectory,源码中会记录 DEBUG 日志 "Picked ... randomly from the configured temporary directories to be used as working directory base.")。

最后再在该基路径下拼上jm_<resourceId>tm_<resourceId>,得到最终的工作目录。

配置 Working Directory

核心配置项

Working Directory 相关配置在 ClusterOptions.java 中定义,共三个配置项,均属于EXPER专家级集群配置(Documentation.Sections.EXPERT_CLUSTER):

配置项作用默认值
process.working-dir所有 Flink 进程共用的工作目录基路径<WORKING_DIR_BASE>未配置时,默认从io.tmp.dirs随机挑选一个目录
process.jobmanager.working-dir仅 JobManager 使用的工作目录基路径未配置时回退到process.working-dir
process.taskmanager.working-dir仅 TaskManager 使用的工作目录基路径未配置时回退到process.working-dir

三点重要说明:

  1. 必须指向本地目录process.working-dir的官方描述是 "Local working directory for Flink processes",它需要指向一个本地目录,而不是分布式文件系统路径。
  2. 专属配置优先,通用配置兜底。源码中JOB_MANAGER_PROCESS_WORKING_DIR_BASETASK_MANAGER_PROCESS_WORKING_DIR_BASE都通过withFallbackKeys(PROCESS_WORKING_DIR_BASE.key())声明了回退键(fallback key),因此进程级配置未设置时会自动读取process.working-dir
  3. 推荐显式配置持久化基路径。默认行为(从io.tmp.dirs随机挑选)意味着每次启动路径都可能变化,若希望进程重启后仍能访问旧的工作目录,就必须显式配置一个稳定的本地基路径。

进程资源 ID 配置

工作目录名中包含的资源 ID 决定了目录的确定性,相关配置项为:

配置项作用默认值源码定义
jobmanager.resource-id指定 JobManager 进程的资源 ID未配置时为随机 UUIDJobManagerOptions.java
taskmanager.resource-id指定 TaskManager 进程的资源 ID未配置时为由 RpcAddress、RpcPort 和 6 位随机字符串组成的随机值TaskManagerOptions.java

源码注释明确说明:

  • JobManager 的jobmanager.resource-idIf not configured, the ResourceID will be generated randomly(随机生成)。
  • TaskManager 的taskmanager.resource-idIf not configured, the ResourceID will be generated with the "RpcAddress:RpcPort" and a 6-character random string. Notice that this option is not valid in Yarn and Native Kubernetes mode.(由 RpcAddress:RpcPort 加 6 位随机字符串组成,且该选项在 Yarn 和 Native Kubernetes 模式下不生效)。

由于随机 ID 会导致每次重启生成不同的工作目录名,若要实现跨重启恢复,就必须为进程显式配置确定性的资源 ID(详见下文"跨进程重启的本地恢复")。

完整配置示例

在 Standalone 部署的conf/flink-conf.yaml中可作如下配置:

# 进程工作目录基路径(本地目录,必须存在且可写) process.working-dir: /path/to/working/dir/base # 可选:JobManager / TaskManager 各自独立的基路径(优先级高于 process.working-dir) # process.jobmanager.working-dir: /path/to/jm/working/dir/base # process.taskmanager.working-dir: /path/to/tm/working/dir/base # 可选:指定确定性资源 ID(跨重启恢复必需) # jobmanager.resource-id: JobManager_1 # taskmanager.resource-id: TaskManager_1

配置完成后,启动 Standalone 集群,日志中会打印所使用的 Working Directory(TaskManager 侧见 TaskManagerRunner.java 的LOG.info("Using working directory: {}", workingDirectory),JobManager 侧见 ClusterEntrypoint.java 的LOG.info("Using working directory: {}.", workingDirectory)),可以直接观察目录是否落在了预期位置。

工作目录中存储的制品

Flink 进程会把以下三类制品写入工作目录:

1. Blob 存储(BlobServer / BlobCache)

JobManager 的 BlobServer 与 TaskManager 的 BlobCache 使用工作目录下的blobStorage/子目录存放分布式缓存、用户 JAR 等 Blob 数据。源码中,JobManager 与 TaskManager 启动时都会把workingDirectory.unwrap().getBlobStorageDirectory()传给 Blob 服务组件(见 ClusterEntrypoint.java 与 TaskManagerRunner.java)。

2. 本地状态(local recovery)

state.backend.local-recovery(新版键名为execution.state-recovery.from-local,见 StateRecoveryOptions.java)开启时,状态后端会把本地快照写入工作目录的localState/子目录。该配置的官方描述强调:

This option configures local recovery for the state backend, which indicates whether to recovery from local snapshot. By default, local recovery is deactivated. Local recovery currently only covers keyed state backends (including both the EmbeddedRocksDBStateBackend and the HashMapStateBackend).

即本地恢复默认关闭,且目前只覆盖 keyed state 后端(EmbeddedRocksDBStateBackend 与 HashMapStateBackend)。TaskManager 侧会把WorkingDirectory整体传给状态后端相关组件(见 TaskManagerRunner.java)。

3. RocksDB 工作目录

若状态后端为 RocksDB,RocksDB 自身的工作目录同样位于进程工作目录之下,从而保证 RocksDB 的本地数据文件在进程重启后仍然可被定位与复用。

此外,从WorkingDirectory源码可以看到,工作目录还包含slotAllocationSnapshots/(槽位分配快照)子目录,供调度相关组件使用。

跨进程重启的本地恢复

工作原理

Working Directory 的核心用途之一,是配合本地恢复特性实现跨进程重启的状态快速恢复(FLIP-201 的设计目标):进程重启后,Flink 可以直接从本地工作目录读取状态快照,无需再从远程存储恢复状态信息,从而显著缩短恢复时间。

要启用这一能力,需要同时满足三个前提条件:

  1. 开启本地恢复:配置state.backend.local-recovery: true
  2. TaskManager 使用确定性资源 ID:通过taskmanager.resource-id显式指定,保证重启前后资源 ID 一致,从而工作目录名(tm_<TM_RESOURCE_ID>)不变;
  3. 失败进程以相同工作目录重启:重启后的 TaskManager 必须能够访问原来的工作目录(同一台机器、同一个本地卷、以相同身份启动)。

配置示例

文档 working_directory.md 给出了最小可用配置:

process.working-dir: /path/to/working/dir/base state.backend.local-recovery: true taskmanager.resource-id: TaskManager_1 # important: Change for every TaskManager process

注意配置中的关键提示:每个 TaskManager 进程都必须使用不同的taskmanager.resource-id。这是因为工作目录名以资源 ID 区分,如果多个 TaskManager 共用一个 ID,它们将写入同一个tm_<ID>目录并互相干扰。

生命周期管理:目录的创建与清理

从源码看,工作目录的创建与清理都遵循"确定性优先"的原则:

  • 创建WorkingDirectory.create(...)在进程启动时确保目录存在并初始化各子目录(见 WorkingDirectory.java)。JobManager 与 TaskManager 分别通过ClusterEntrypointUtils.createJobManagerWorkingDirectory/createTaskManagerWorkingDirectory完成(ClusterEntrypointUtils.java)。
  • 清理:进程正常结束时(或工作目录非确定性时)会删除整个工作目录。TaskManager 侧逻辑见 TaskManagerRunner.java:if (!workingDirectory.isDeterministic() || terminationResult == Result.SUCCESS) { workingDirectory.unwrap().delete(); }。JobManager 侧逻辑与之对称(ClusterEntrypoint.java)。

这段逻辑的含义是:当资源 ID 为随机生成(工作目录不确定,isDeterministic() == false)时,无论进程因何退出都会清理工作目录;当资源 ID 确定时,仅在进程正常退出(Result.SUCCESS)时才清理,异常失败则保留目录——这正是跨重启恢复得以成立的关键:失败重启后目录仍然存在,本地状态得以复用。

配套测试验证

仓库中配套的单元测试对上述行为有直接验证:

  • WorkingDirectoryTest.java:验证WorkingDirectory创建时目录结构(tmp/localState/blobStorage/slotAllocationSnapshots/)是否正确生成。
  • ClusterEntrypointTest.java:覆盖工作目录生成与生命周期相关行为。

最佳实践与注意事项

何时必须显式配置 Working Directory

默认情况下 Flink 会从io.tmp.dirs随机挑选基路径,这对仅需临时目录的普通运行没有问题。但以下场景必须显式配置

  • 期望 TaskManager 失败重启后能本地恢复状态(配合state.backend.local-recovery);
  • 需要把 Blob、本地状态等数据放置到特定的高性能本地磁盘(如 NVMe SSD),而不是默认临时目录;
  • 需要多个进程工作目录相互隔离、可预期(例如监控、排障时能快速定位某个进程的目录)。

配置检查清单

检查项说明
基路径是本地目录不能指向 HDFS、S3 等分布式存储;目录需要存在且对运行用户可写
state.backend.local-recovery仅覆盖 keyed stateRocksDB 与 HashMap 状态后端支持,其他状态后端不适用
每个 TaskManager 使用独立taskmanager.resource-id多个 TM 共用 ID 会写入同一目录并互相干扰
重启保持相同身份与存储卷访问进程用户、挂载的本地卷需与之前一致,否则无法读取旧目录
Yarn / Native Kubernetes 模式下taskmanager.resource-id不生效该配置项在上述资源提供者模式下无效,请使用对应模式下的资源 ID 管理机制

io.tmp.dirs的关系

当不配置任何 working-dir 选项时,基路径从io.tmp.dirs中随机选取(见 ClusterEntrypointUtils.java)。这意味着io.tmp.dirs中每个候选目录都可能成为工作目录基路径。在生产环境中,建议将工作目录与临时目录分离管理:显式配置process.working-dir指向持久化本地卷,让io.tmp.dirs继续承担纯粹的临时文件职责。

小结

Working Directory 是 Flink 进程本地持久化运行时信息的核心机制:它以<WORKING_DIR_BASE>/jm_<JM_RESOURCE_ID><WORKING_DIR_BASE>/tm_<TM_RESOURCE_ID>的规则组织目录,统一承载 Blob 存储、本地状态快照与 RocksDB 工作目录,并借助确定性资源 ID + 失败不清除目录的生命周期策略,支撑起跨进程重启的本地快速恢复。配置层面只需掌握三条主线:基路径选项(process.working-dir及其进程级变体)、资源 ID 选项(jobmanager.resource-id/taskmanager.resource-id)以及本地恢复开关(state.backend.local-recovery)。理解并正确配置这三个维度,即可在 Standalone 集群中实现进程重启后的本地状态复用,显著降低故障恢复成本。

  • 大数据
  • 流处理
  • 批处理
  • 数据工程

【免费下载链接】flink

项目地址:https://gitcode.com/gh_mirrors/fli/flink
点击查看免费下载

相关推荐

创作声明:本文部分内容由AI辅助生成(AIGC),仅供参考

版权声明: 本文来自互联网用户投稿,该文观点仅代表作者本人,不代表本站立场。本站仅提供信息存储空间服务,不拥有所有权,不承担相关法律责任。如若内容造成侵权/违法违规/事实不符,请联系邮箱:809451989@qq.com进行投诉反馈,一经查实,立即删除!
网站建设 2026/9/24 7:38:39

基于RK3588的8K全景相机:多路采集、NPU拼接与8K编码全链路实践

/* MD / 富文本中的 .toc(含博客园搬家等嵌套结构);.toc-box 在侧栏,不受影响 */#content_views .toc,/* 编辑器常在目录前后插入空 p(:empty 仍占 20px),一并去掉避免顶空隙 */#content_views.markdown_views > p:empty:has(+ .toc),#content_views.markdown_views …

作者头像 李华
网站建设 2026/9/24 7:31:32

ESP32-S3-BOX-3实战:智能语音与物联网联动开发指南

/* MD / 富文本中的 .toc(含博客园搬家等嵌套结构);.toc-box 在侧栏,不受影响 */#content_views .toc,/* 编辑器常在目录前后插入空 p(:empty 仍占 20px),一并去掉避免顶空隙 */#content_views.markdown_views > p:empty:has(+ .toc),#content_views.markdown_views …

作者头像 李华
网站建设 2026/9/24 7:14:31

Modbus转MQTT数据采集全流程:从RS485到云端实战指南

/* MD / 富文本中的 .toc(含博客园搬家等嵌套结构);.toc-box 在侧栏,不受影响 */#content_views .toc,/* 编辑器常在目录前后插入空 p(:empty 仍占 20px),一并去掉避免顶空隙 */#content_views.markdown_views > p:empty:has(+ .toc),#content_views.markdown_views …

作者头像 李华