首页 / 资讯中心 / 文章详情

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

Flink Working Directory 完全指南:进程工作目录配置与跨重启本地恢复实战 ★ FEATURED ARTICLE
大数据流处理批处理数据工程【免费下载链接】flink项目地址https://gitcode.com/gh_mirrors/fli/flink点击查看免费下载Working Directory工作目录是 Flink 为 JobManager 与 TaskManager 进程提供的本地持久化目录用于存放进程重启后可以恢复的运行时信息其设计由 FLIP-198 引入、并由 FLIP-201 扩展出跨进程重启的本地恢复能力。本文以当前仓库Apache FlinkStandalone 部署文档 working_directory.md 为主体结合flink-runtime/flink-core的源码实现系统讲解 Working Directory 的目录结构、配置项、进程资源 ID 的确定规则、存储在其中的三类制品以及如何利用它实现 TaskManager 重启后从本地快速恢复状态。读完本文你将能够为 Standalone 集群正确配置工作目录并搭建一套进程重启不丢失本地状态的高可用实践方案。什么是 Working DirectoryWorking 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.cleanDirectorylocalState/本地状态目录供本地恢复local recovery使用blobStorage/Blob 存储目录供 BlobServer / BlobCache 使用slotAllocationSnapshots/槽位分配快照目录这些子目录可以通过WorkingDirectory提供的getTmpDirectory()、getLocalStateDirectory()、getBlobStorageDirectory()、getSlotAllocationSnapshotDirectory()等方法在运行时获取。工作目录的目录结构Flink 为两类进程分别规划了工作目录其命名规则如下JobManager 工作目录WORKING_DIR_BASE/jm_JM_RESOURCE_IDTaskManager 工作目录WORKING_DIR_BASE/tm_TM_RESOURCE_ID其中WORKING_DIR_BASE是工作目录基路径baseJM_RESOURCE_ID是 JobManager 进程的资源 IDresource idTM_RESOURCE_ID是 TaskManager 进程的资源 ID。也就是说同一个基路径下不同进程通过jm_/tm_前缀与各自的资源 ID 区分目录。这一规则与源码中目录生成逻辑完全一致在 ClusterEntrypointUtils.java 中generateTaskManagerWorkingDirectoryFile使用tm_ resourceId作为目录名generateJobManagerWorkingDirectoryFile使用jm_ resourceId作为目录名。整个生成流程generateWorkingDirectoryFile见 ClusterEntrypointUtils.java的决策顺序是若配置了进程专属的 working-dir 选项如process.jobmanager.working-dir/process.taskmanager.working-dir则直接以其为基路径否则若配置了通用选项process.working-dir则以其为基路径专属选项通过withFallbackKeys回退到通用选项若以上均未配置则从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-dirprocess.taskmanager.working-dir仅 TaskManager 使用的工作目录基路径未配置时回退到process.working-dir三点重要说明必须指向本地目录。process.working-dir的官方描述是 Local working directory for Flink processes它需要指向一个本地目录而不是分布式文件系统路径。专属配置优先通用配置兜底。源码中JOB_MANAGER_PROCESS_WORKING_DIR_BASE和TASK_MANAGER_PROCESS_WORKING_DIR_BASE都通过withFallbackKeys(PROCESS_WORKING_DIR_BASE.key())声明了回退键fallback key因此进程级配置未设置时会自动读取process.working-dir。推荐显式配置持久化基路径。默认行为从io.tmp.dirs随机挑选意味着每次启动路径都可能变化若希望进程重启后仍能访问旧的工作目录就必须显式配置一个稳定的本地基路径。进程资源 ID 配置工作目录名中包含的资源 ID 决定了目录的确定性相关配置项为配置项作用默认值源码定义jobmanager.resource-id指定 JobManager 进程的资源 ID未配置时为随机 UUIDJobManagerOptions.javataskmanager.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 DirectoryTaskManager 侧见 TaskManagerRunner.java 的LOG.info(Using working directory: {}, workingDirectory)JobManager 侧见 ClusterEntrypoint.java 的LOG.info(Using working directory: {}., workingDirectory)可以直接观察目录是否落在了预期位置。工作目录中存储的制品Flink 进程会把以下三类制品写入工作目录1. Blob 存储BlobServer / BlobCacheJobManager 的 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 工作目录若状态后端为 RocksDBRocksDB 自身的工作目录同样位于进程工作目录之下从而保证 RocksDB 的本地数据文件在进程重启后仍然可被定位与复用。此外从WorkingDirectory源码可以看到工作目录还包含slotAllocationSnapshots/槽位分配快照子目录供调度相关组件使用。跨进程重启的本地恢复工作原理Working Directory 的核心用途之一是配合本地恢复特性实现跨进程重启的状态快速恢复FLIP-201 的设计目标进程重启后Flink 可以直接从本地工作目录读取状态快照无需再从远程存储恢复状态信息从而显著缩短恢复时间。要启用这一能力需要同时满足三个前提条件开启本地恢复配置state.backend.local-recovery: trueTaskManager 使用确定性资源 ID通过taskmanager.resource-id显式指定保证重启前后资源 ID 一致从而工作目录名tm_TM_RESOURCE_ID不变失败进程以相同工作目录重启重启后的 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.javaif (!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点击查看免费下载相关推荐Flink Standalone 部署的 Working Directory 配置完全指南原理、配置项与跨进程重启的本地恢复Flink Standalone 部署的 Working Directory 配置完全指南原理、配置项与跨进程重启的本地恢复 本指南以 Flink 官方文档大数据流处理批处理数据工程Hydra 工作目录Working Directory自定义完全指南run/sweep 输出目录模式详解Hydra 工作目录Working Directory自定义完全指南run/sweep 输出目录模式详解 导读 Hydra 会自动为每次运行run和多开发工具后端CLIWezTerm 中 wezterm set-working-directory 命令完全指南OSC 7 工作目录通知的原理、用法与集成实践WezTerm 中 wezterm set working directory 命令完全指南OSC 7 工作目录通知的原理、用法与集成实践 导读 wezter桌面应用开发工具跨平台创作声明:本文部分内容由AI辅助生成(AIGC),仅供参考
阅读完成 · 觉得有帮助?
咨询建站