Ray Tune 同步配置(SyncConfig)完全指南:实验状态、Checkpoint 与 Artifact 的持久化同步
【免费下载链接】rayRay is an AI compute engine. Ray consists of a core distributed runtime and a set of AI Libraries for accelerating ML workloads.项目地址: https://gitcode.com/gh_mirrors/ra/ray
导读
本文围绕 Ray Tune 的同步(Syncing)机制展开,核心讲解ray.tune.SyncConfig这一配置对象:在分布式超参数调优中,它控制实验目录(含搜索器状态、trial 列表与元数据)如何从 driver 同步到持久化存储(如 AWS S3、GCS、NFS),以及 trial 目录中的任意产物(artifacts)是否随tune.report上传到云端。读完本文,你将掌握SyncConfig全部四个参数(sync_period、sync_timeout、sync_artifacts、sync_artifacts_on_checkpoint)的语义与默认值,能够为单机、NFS、云存储三种场景配置持久化存储,理解同步任务在源码层(Syncer抽象基类与后台进程)如何被节流、超时与重试,并知道如何从存储 URI 恢复被中断的实验。
一、Tune 中同步发生在哪里
根据 python/ray/tune/syncer.py 中SyncConfig的类文档,Ray Tune 的文件同步(主要是上传)发生在两个层面:
- 实验驱动端(head node):实验 driver 将整个实验目录同步到存储(
RunConfig(storage_path)指定的位置),其中包含实验状态——搜索器(searcher)状态、trial 列表及其状态、trial 元数据等。这是实验级容错与恢复的基础。 - trial 端:通过设置
sync_artifacts=True,可以把 trial 目录中的产物(你在训练过程中随手 dump 到ray.tune.get_context().get_trial_dir()下的任意文件)同步到存储。一个包含大量 trial 的 Tune 实验中,每个 trial 都会把自己的 trial 目录上传到存储。
这两类数据对应了 Ray Tune 持久化存储要解决的四个核心场景(见配套用户指南 doc/source/tune/tutorials/tune-storage.rst):
- Trial 级容错:trial 恢复(如节点故障、实验暂停后重启)时可能被调度到不同节点,但必须能访问其最新 checkpoint;
- 实验级容错:整个实验恢复(如集群意外崩溃)时,Tune 需要访问最新实验状态与全部 trial checkpoint,从断点继续;
- 实验后分析:集群终止后,仍有一个集中位置存放所有 trial 的数据,便于分析最优 checkpoint 与超参配置;
- 下游任务衔接:用配置好的存储把模型与产物交付给推理/批处理等下游任务。
SyncConfig正是控制上述同步行为(频率、超时、是否同步产物)的配置对象,属于@PublicAPI(stability="beta")的公开 API。
二、SyncConfig 参数详解
SyncConfig在 python/ray/tune/syncer.py 中定义,是一个继承自ray.train._internal.syncer.SyncConfig的 dataclass。基础字段定义于 python/ray/train/_internal/syncer.py(DEFAULT_SYNC_PERIOD = 300,DEFAULT_SYNC_TIMEOUT = 1800),Tune 侧在此基础上增加了 artifact 同步开关。全部参数如下:
| 参数 | 类型 | 默认值 | 说明 |
|---|---|---|---|
sync_period | int | 300(5 分钟) | 两次同步操作之间的最小等待秒数。值越小,存储中的数据更新越频繁,但同步开销越大。 |
sync_timeout | int | 1800(30 分钟) | 单个同步进程允许运行的最大秒数。同步操作运行超过该时长会抛出TimeoutError。 |
sync_artifacts | bool | False | [Beta] 是否把保存到 trial 目录(通过ray.tune.get_context().get_trial_dir()访问)的产物同步到tune.RunConfig(storage_path)配置的持久化存储。trial 或远程 worker 会在每次tune.report时尝试发起一次产物同步(受sync_period与sync_artifacts_on_checkpoint约束)。默认False,即默认不持久化任何产物。 |
sync_artifacts_on_checkpoint | bool | True | 若为True,在每次上报的 checkpoint 上强制同步 trial/worker 产物。仅在sync_artifants为True时生效。 |
参数语义要点:
sync_period是"节流器":它规定两次同步之间最少间隔多久,防止高频率tune.report触发无节制的网络传输。源码中该节流逻辑实现在Syncer.sync_up_if_needed/sync_down_if_needed(见 python/ray/train/_internal/syncer.py),只有满足now - self.last_sync_up_time >= self.sync_period时才会真正发起sync_up,并更新last_sync_up_time。sync_timeout是"看门狗":_BackgroundProcess.wait会从进程启动时刻起计时,超过timeout秒仍存活则清空进程并抛出TimeoutError(错误信息形如xxx did not finish running within the timeout of N seconds.)。sync_artifacts默认关闭:因为产物是"任意文件",可能包含大体积中间结果,默认不同步以免拖慢训练;需要持久化产物时才显式开启。sync_artifacts_on_checkpoint默认开启:只要开启了产物同步,就会在每个 checkpoint 上报时强制同步一次。
SyncConfig也可通过tune.TuneConfig(sync_config=...)或tune.RunConfig传入。在 python/ray/tune/tune.py 的run入口中可以看到sync_config = sync_config or SyncConfig()——未显式提供时使用默认配置;Tuner内部则在 python/ray/tune/impl/tuner_internal.py 中以sync_config=self._run_config.sync_config将其注入存储上下文。
三、为 Tune 配置持久化存储的三种方式
SyncConfig与RunConfig(storage_path)配合使用。根据 doc/source/tune/tutorials/tune-storage.rst,Tune 支持三种存储场景。
3.1 云存储(AWS S3、Google Cloud Storage)
当集群所有节点都能访问云存储时,把所有实验输出保存到共享 bucket。只需让 Ray Tune上传到远程storage_path:
from ray import tune tuner = tune.Tuner( trainable, run_config=tune.RunConfig( name="experiment_name", storage_path="s3://bucket-name/sub-path/", ) ) tuner.fit()运行后所有实验结果位于s3://bucket-name/sub-path/experiment_name。注意:
- head 节点本地并不会保留全部实验结果;如需进一步处理最优 checkpoint,需要先从云存储拉取;
- 恢复实验时应使用云存储 URI 对应的实验目录,而非 head 节点上的本地实验目录(见下文恢复示例)。
3.2 网络文件系统(NFS)
当所有 Ray 节点都能访问 NFS(如 AWS EFS、Google Cloud Filestore)时,把共享文件系统路径直接作为结果保存路径即可:
from ray import tune tuner = tune.Tuner( trainable, run_config=tune.RunConfig( name="experiment_name", storage_path="/mnt/path/to/shared/storage/", ) ) tuner.fit()运行后所有实验结果位于/path/to/shared/storage/experiment_name。NFS 属于"所有节点共享同一路径",本质上无需跨节点复制,是成本最低的共享方案。
3.3 单节点本地文件系统
在单节点(如笔记本)上运行实验时,Tune 默认使用本地文件系统作为 checkpoint 与其他产物的存储位置:结果默认保存到~/ray_results下带唯一自动生成名称的子目录,除非通过RunConfig的storage_path和name自定义:
from ray import tune tuner = tune.Tuner( trainable, run_config=tune.RunConfig( storage_path="/tmp/custom/storage/path", name="experiment_name", ) ) tuner.fit()运行后实验结果位于/tmp/custom/storage/path/experiment_name。NFS 或云存储也可用于单机实验——例如实例终止会清空本地存储时,把结果持久化到外部存储是很有用的做法。
需要特别强调的是:在多节点集群上使用 head 节点的本地文件系统作为持久化存储已被弃用(Deprecated)。如果保存 trial checkpoint 且运行在多节点集群,未配置 NFS 或云存储时 Tune 默认会直接报错。
四、实战示例:云端实验与中断恢复
下面给出带SyncConfig与 checkpoint 保留策略的完整示例(源于 doc/source/tune/tutorials/tune-storage.rst 的实战骨架,假设在集群 head 节点运行,my_trainable实现了 checkpoint 的保存与加载):
import os import ray from ray import tune from your_module import my_trainable tuner = tune.Tuner( my_trainable, run_config=tune.RunConfig( # 实验名称 name="my-tune-exp", # 配置实验数据与 checkpoint 的持久化方式。 # 推荐云存储 checkpoint:集群实例被终止后数据依然存在,且性能更好。 storage_path="s3://my-checkpoints-bucket/path/", checkpoint_config=tune.CheckpointConfig( # 始终保留最优的 5 个 checkpoint # (按 trainable 上报的 max-auc 指标排序) checkpoint_score_attribute="max-auc", checkpoint_score_order="max", num_to_keep=5, ), # 显式控制同步行为(可选) sync_config=tune.SyncConfig( sync_period=120, # 每 120 秒至多同步一次 sync_timeout=600, # 单次同步最多运行 600 秒 sync_artifacts=True, # 同步 trial 目录中的任意产物 ), ), ) # 启动运行 results = tuner.fit()该例中 trial checkpoint 保存到s3://my-checkpoints-bucket/path/my-tune-exp/<trial_name>/checkpoint_<step>。
若运行因任何原因中断(用户 CTRL+C、OOM 被终止等),可随时基于云端的实验状态恢复:
from ray import tune tuner = tune.Tuner.restore( "s3://my-checkpoints-bucket/path/my-tune-exp", trainable=my_trainable, resume_errored=True, ) tuner.fit()恢复选项包括resume_unfinished、resume_errored与restart_errored,详见Tuner.restore的 API 文档。恢复时必须使用云端实验目录 URI(而非本地目录),这正是因为同步机制保证云端持有最新的实验状态快照。
五、源码视角:同步机制如何工作
5.1 Syncer 抽象基类与命令原语
同步的核心抽象位于 python/ray/train/_internal/syncer.py 的Syncer(@DeveloperAPI,抽象基类),它定义了三类同步方向:
sync_up(local_dir, remote_dir, exclude):把本地目录同步到远程 URI(protocol://remote/path),exclude支持排除模式(如["*/checkpoint_*"]跳过 trial checkpoint);sync_down(remote_dir, local_dir, exclude):反向拉取,用于恢复场景;delete(remote_dir):删除远端存储目录。
同步任务通常是异步的:sync_up/down返回True表示已派发后台进程,之后可调用wait()等待完成(超过sync_timeout抛TimeoutError)。基类还实现了wait_or_retry(max_retries=2, backoff_s=5):等待失败后按 5 秒退避重试最多 2 次,全部失败则抛出包含完整 traceback 的RuntimeError。
5.2 节流、后台进程与超时
_BackgroundSyncer与_BackgroundProcess是默认实现的关键:
- 节流:
_should_continue_existing_sync()检查"上一次同步是否仍在运行且未超过sync_timeout",若是则跳过本次同步(记 debug/warning 日志),避免并发堆积; - 后台进程:
_BackgroundProcess用 daemon 线程包装同步命令,start()记录start_time,wait(timeout)从启动时刻起计时,超时即清空进程并抛TimeoutError; - 失败重试:
_launch_sync_process会先wait()等上一次同步结束(失败则记录 warning),再启动新进程;retry()基于_current_cmd重新创建后台进程执行同一条命令。
也就是说,sync_period、sync_timeout两个参数最终会注入Syncer实例(见Syncer.__init__的sync_period: float = DEFAULT_SYNC_PERIOD、sync_timeout: float = DEFAULT_SYNC_TIMEOUT),从底层控制每次同步的触发时机与存活上限。
5.3 Artifact 同步的触发点
sync_artifacts与sync_artifacts_on_checkpoint的落地位置在 python/ray/tune/trainable/trainable.py:trial/worker 在 checkpoint 上报时调用self._storage.sync_config.sync_artifacts_on_checkpoint决定是否以force=True方式同步产物;普通场景下则每次tune.report时尝试发起一次产物同步,受sync_period节流。
5.4 测试用例佐证
仓库测试验证了上述行为:
- python/ray/tune/tests/test_tuner.py 中以
tune.SyncConfig(sync_artifacts=True)构造配置,验证开启产物同步时的端到端行为; - python/ray/tune/tests/execution/test_controller_checkpointing_integration.py 使用
ray.tune.SyncConfig(sync_timeout=0.5)构造 mock 存储上下文,验证极短超时下的同步控制逻辑。
六、配置建议与注意事项
sync_period权衡:值越小存储中数据越新鲜(故障时丢失窗口越小),但同步开销越大;默认 5 分钟适合大多数训练节奏。若 checkpoint 很大且频率高,可适当调大以降低网络与对象存储成本。sync_timeout权衡:值越大,超大目录(如大模型 checkpoint)越可能完成上传;但卡死的同步进程会占用更久。默认 30 分钟;大 checkpoint 场景建议调大,小文件场景可调小以更快暴露故障。- 产物同步默认关闭:
sync_artifacts=False意味着 trial 目录中的任意文件默认不会持久化;需要保存日志、中间结果或自定义产物时再开启,并注意它受sync_period节流、在每次 checkpoint 时强制同步一次。 - 多节点必须配共享存储:多节点集群中若未配置 NFS 或云存储,保存 trial checkpoint 会直接报错(本地文件系统方案已弃用),请优先选择云存储 checkpoint。
- 恢复用 URI 而非本地路径:中断恢复时使用
Tuner.restore并传入云存储/共享存储中的实验目录 URI。
总结
ray.tune.SyncConfig是 Ray Tune 持久化与同步行为的统一配置入口:sync_period控制同步频率、sync_timeout控制单次同步的存活上限、sync_artifacts/sync_artifacts_on_checkpoint控制 trial 产物的持久化策略。它配合RunConfig(storage_path)即可在单机、NFS、云存储三种场景下获得 trial 级与实验级容错能力。源码层面,Syncer抽象基类通过sync_up/sync_down/delete原语、sync_up_if_needed节流、_BackgroundProcess超时控制与wait_or_retry重试机制,把配置参数落地为稳定可靠的分布式文件同步管线。
【免费下载链接】rayRay is an AI compute engine. Ray consists of a core distributed runtime and a set of AI Libraries for accelerating ML workloads.项目地址: https://gitcode.com/gh_mirrors/ra/ray
创作声明:本文部分内容由AI辅助生成(AIGC),仅供参考