首页
/ ColossalAI ColossalChat 基于 Ray 的 Detached PPO 分布式训练:Maker 与 Trainer 异步解耦实战指南

ColossalAI ColossalChat 基于 Ray 的 Detached PPO 分布式训练:Maker 与 Trainer 异步解耦实战指南

2026-09-05 11:25:26作者:胡唯隽

本文围绕 ColossalAI 仓库中 ColossalChat 的 Ray 模块文档 展开,系统讲解"将经验生成(Experience Maker)与模型训练(Trainer)完全分离"的 Detached PPO 架构:如何用 Ray Actor 编排 Maker 与 Trainer、如何配置分布式环境变量与自定义传输拓扑、以及底层源码中经验回传、模型参数同步与并发组(concurrency groups)的实现机制。读完本文,你可以理解并复现 ColossalChat Stage 3 中推理与训练异步流水线的完整搭建流程。

1. 架构总览:为什么要 Detach Experience Makers 和 Trainers

传统 RLHF(PPO)训练中,经验生成(采样 + 打分)与参数更新往往串行执行在同组进程里,GPU 在"推理等待训练、训练等待推理"之间反复空转。Ray 模块文档提出的 Detached 架构把二者彻底拆分为两类独立角色:

  • Experience Maker:执行推理、产出经验(Experience),并远程投递给 Trainer(文档中的步骤 1);
  • Trainer:消费经验训练模型,并周期性地把新模型参数回传给 Maker(文档中的步骤 2.1、2.2);
  • Experience Buffer:作为中间缓冲,让"参数传输"与"计算"相互重叠(overlap)。

在这种模式下,每个节点可以持续工作、没有模型空闲时间;推理侧与训练侧还能各自采用不同的优化策略——推理侧追求速度(甚至可以做量化),训练侧追求显存效率(Zero/Gemini),这同样有利于横向扩展。

文档明确指出,DetachedPPOTrainerExperienceMakerHolderRay Actor(注意文档特意强调这是 Ray 的 Actor 概念,区别于强化学习中的 Actor Model),分别对应上图中的 Trainer 与 Experience Maker。文档建议读者参考 Ray 官方 Core 文档来理解 Actor 与命名 Actor(ray.get_actor(name))机制。

1.1 角色对应的源码位置

文档角色 源码实现 说明
Experience Maker 宿主 experience_maker_holder.py ExperienceMakerHolder,Ray 远程 Actor
PPO Trainer detached_trainer_ppo.py DetachedPPOTrainer,Ray 远程 Actor
Trainer 抽象基类 detached_trainer_base.py DetachedTrainer,封装 buffer 与训练循环
工具函数 utils.py 模型/策略/Tokenizer 工厂、环境变量注入

2. 核心组件实现解析

2.1 ExperienceMakerHolder:推理侧的 Ray Actor

ExperienceMakerHolder 的类装饰器(experience_maker_holder.py#L21-L22)声明了三个并发组:

@ray.remote(concurrency_groups={"experience_io": 1, "model_io": 1, "compute": 1})
class ExperienceMakerHolder:

这三个并发组把 Maker 内部三条互不干扰的路径隔离开:compute 做经验生成(_make_experience),experience_io 把经验发送出去(_send_items),model_io 接收 Trainer 回传的参数(update_experience_maker)。因此 Trainer 推送参数时不会阻塞推理循环,这正是文档所说"用 buffer 重叠传输与计算"的代码基础。

构造参数(见 experience_maker_holder.py#L31-L46)包括:

参数 类型 / 默认值 含义
detached_trainer_name_list List[str] 目标 Trainer 的 Ray Actor 名列表,经验将投递给它们
strategy_fn Callable[[], Strategy] 返回训练/推理策略对象(DDP、Zero、Gemini 等)的工厂
model_fn Callable 返回 (actor, critic, reward_model, initial_model) 四元组
env_info Dict[str, str] 分布式环境变量(rank/world_size/master_addr 等)
sync_models_from_trainers bool,默认 False 若为 True,Maker 需等待 Trainer 调用 sync_models_to_remote_makers() 完成初始化同步
buffer_cpu_offload bool,默认 True 经验发送前是否先 offload 到 CPU 以释放 GPU 显存
kl_coef float,默认 0.1 KL 散度奖励惩罚系数
callbacks List[MakerCallback] 训练回调钩子
update_lora_weights bool,默认 False 只同步 LoRA 增量权重而非全量参数

经验产出与发送的核心路径在 _inference_stepexperience_maker_holder.py#L128-L139):加锁生成经验 → 可选地 experience.to_device("cpu") → 调用 _send_items_send_items 会把一个批次的经验用 split_experience_batch 拆成单条 BufferItem,然后按 _target_idx **轮询(round-robin)**分配给 target_trainer_list 中的每个 Trainer,并通过 target_trainer.buffer_extend.remote(items) 异步投递(experience_maker_holder.py#L116-L126)。这意味着多个 Trainer 时经验是交错分发的,天然均衡了各 Trainer 的训练负载。

参数回传接口 update_experience_makerexperience_maker_holder.py#L170-L229)采用分块(chunked)协议

  1. Trainer 先发 chunk_start=True 标记,Maker 侧 acquire() 模型访问锁——此刻推理暂停;
  2. Trainer 循环发送 new_actor_state_dict / new_critic_state_dict 的各个分片,Maker 逐片 load_state_dict(..., strict=False)
  3. Trainer 最后发 chunk_end=True,Maker 释放锁,恢复推理;若 fully_update=True,还置 _is_fully_initialized = True,解除 Maker 启动时的初始化等待。

若启用 update_lora_weights,Maker 会通过 LoRAConstructor 只重建 LoRA 增量(reconstruct_increase / load_state_dict_increase),大幅减少回传流量。

2.2 DetachedPPOTrainer:训练侧的 Ray Actor

DetachedPPOTrainer 的并发组更多(detached_trainer_ppo.py#L19-L21):

@ray.remote(
    concurrency_groups={"buffer_length": 1, "buffer_append": 1, "buffer_sample": 1, "model_io": 1, "compute": 1}
)

其中 buffer_length / buffer_append / buffer_sample 分别服务于 buffer_get_length / buffer_extend / _buffer_sample 三个由 Maker 或自身调用的缓冲接口,model_io 承载参数回传 _update_remote_makerscompute 承载 training_step。同样地,向 Maker 推送参数(model_io 组)与本地训练(compute 组)可以并行。

构造参数(见 detached_trainer_ppo.py#L43-L58):

参数 默认值 含义
experience_maker_holder_name_list 必填 目标 Maker 的 Ray Actor 名列表,训练后参数回传对象
strategy_fn 必填 训练策略工厂
model_fn 必填 返回 (actor, critic) 二元组的工厂
train_batch_size 8 训练批大小
buffer_limit 0 replay buffer 上限(0 表示不限)
eps_clip 0.2 PPO 策略损失裁剪系数(PolicyLoss
value_clip 0.4 价值损失裁剪系数(ValueLoss
dataloader_pin_memory True 训练 DataLoader 是否 pin memory
update_lora_weights False 只向 Maker 同步 LoRA 增量

注意一个实现细节:当策略是 LowLevelZeroStrategyGeminiStrategy 时,Trainer 自动选用 ColossalAI 的 HybridAdam 优化器,否则用 torch.optim.Adamdetached_trainer_ppo.py#L74-L79)——这就是文档"不同优化策略可分别应用于推理与训练"的具体体现之一。

训练循环在基类 DetachedTrainer.fit 中(detached_trainer_base.py#L105-L114):外层按 total_steps // update_steps 划分 episode,每个 episode 内 _learn(update_steps, train_epochs) 消费 buffer,随后调用 _update_remote_makers() 向目标 Maker 分块推送新参数。_learn 的首个 epoch 是"热启动":若 buffer 尚未被 Maker 填满,Trainer 会自旋轮询 buffer_get_length 等待经验到达,而不是提前训练陈旧数据。

_update_remote_makers 的回传顺序保证了正确性(detached_trainer_ppo.py#L102-L139):先对所有目标 Maker ray.get 一次 chunk_start=True(确保全部锁住后才发参数),然后分片发送 actor 参数、再分片发送 critic 参数,最后广播 chunk_end=True。若 fully_update=False,还只取需要梯度的参数(requires_grad_only=True)。

2.3 自定义 Maker↔Trainer 传输拓扑

文档强调:通过 .options(name="str") 给 Trainer 与 Maker 起同名风格的命名 Actor(trainer_i / maker_i),再分别用 detached_trainer_name_list(Maker 指向哪些 Trainer)和 experience_maker_holder_name_list(Trainer 指向哪些 Maker)即可自定义整个传输图

命名解析依赖 os.environ["RAY_NAMESPACE"],例如 experience_maker_holder.py#L100-L104 中的 ray.get_actor(name, namespace=...)utils.py 还提供了 get_receivers_per_senderutils.py#L118-L130),用"取模轮询"规则计算某个发送者应该把数据发给哪些接收者,便于按 2 Maker : 1 Trainer、2 Maker : 2 Trainer 等比例自动生成拓扑。

3. Usage:从文档到代码的完整搭建流程

文档正文给出的示例入口是 ColossalAI/application/Chat/examples/ray(在当前仓库中已迁移为 examples/community/ray 目录)。以下按文档的三个章节逐步展开,并修正了原文档中个别笔误(如 return = torch.utils.data.DataLoader(...)),保证代码可直接理解与运行。

3.1 前置:准备 Ray 集群与提交任务

community/ray 示例说明,标准流程是:

  1. 按 Ray 官方指南搭建支持 GPU 的 Ray 集群,记录 API server 端点(形如 http://your.head.node.address:8265);
  2. 克隆本仓库;
  3. 通过 Job Submission Client 提交任务:
python applications/ColossalChat/examples/community/ray/ray_job_script.py http://your.head.node.address:8265
  1. 打开 Ray Dashboard 查看训练作业。

ray_job_script.py 内部就是一个最小示例:

client = JobSubmissionClient(api_server_endpoint)
client.submit_job(
    entrypoint="python .../train_prompts_on_ray.py --strategy colossalai_zero2 --prompt_csv_url <csv_url>",
    runtime_env={"working_dir": "applications/Chat", "pip": ["torch==1.13.1", ..., "colossalai>=0.2.4"]},
)

train_prompts_on_ray.py 展示了一个"全解耦"的 PPO 实现:Actor、Critic、Initial Model、Reward Model 各自是一组 @ray.remote(num_gpus=1) 的 Ray Actor,每节点内通过 MASTER_ADDR/MASTER_PORT/WORLD_SIZE/RANK/LOCAL_RANK 环境变量自行初始化 torch.distributed 进程组(见 DistributedTorchRayActor 类),多 GPU 时每类模型用 placement_group(strategy="STRICT_SPREAD") 锁定资源。其 CLI 参数包括 --strategy {ddp,colossalai_gemini,colossalai_zero2}--model {gpt2,bloom,opt}--num_actor_nodes/--num_critic_nodes/--num_initial_nodes/--num_reward_nodes/--num_gpus_per_node 等,直观体现了"每种 PPO 角色可独立指定节点数"的弹性部署思想。

注意:文档开头带有警告标记,说明该内容在 ColossalChat 大版本更新后可能滞后。以当前仓库源码为准:示例脚本属于早期风格(各角色为独立 Ray Actor 类),而 coati/ray 模块本身已演进为 ExperienceMakerHolder / DetachedPPOTrainer 的统一封装,本文后续步骤即基于该封装。

3.2 Setup Makers(配置经验生成端)

第一步:定义 Maker 的分布式环境变量。 每个 Maker rank 需要一份 env_infoutils.py 的 set_dist_env 会把它们写入 RANK/LOCAL_RANK/WORLD_SIZE/MASTER_PORT/MASTER_ADDR):

env_info_makers = [{
    'local_rank': '0',
    'rank': str(rank),
    'world_size': str(num_makers),
    'master_port': maker_port,
    'master_addr': master_addr
} for rank in range(num_makers)]

第二步:定义 Maker 侧模型。 Maker 需要四个模型:策略网络 actor、价值网络 critic、打分网络 reward_model、以及冻结的初始策略 initial_model(用于计算 KL 惩罚)。utils.py 提供了对应工厂函数(支持 gpt2 / bloom / opt / llama):

from coati.ray.utils import get_actor_from_args, get_critic_from_args, get_reward_model_from_args

def model_fn():
    actor = get_actor_from_args(model, pretrained, lora_rank=0)
    critic = get_critic_from_args(model, pretrained)
    reward_model = get_reward_model_from_args(model, pretrained)
    initial_model = get_actor_from_args(model, pretrained)
    return actor, critic, reward_model, initial_model

第三步:创建 Maker 命名 Actor。

experience_holder_refs = [
    ExperienceMakerHolder.options(
        name=f"maker_{i}",
        num_gpus=1,
        max_concurrency=2
    ).remote(
        detached_trainer_name_list=[f"trainer_{x}" for x in target_trainers(i, num_makers, num_trainers)],
        strategy_fn=lambda: get_strategy_from_args("colossalai_zero2"),
        model_fn=model_fn,
        env_info=env_info_makers[i],
        sync_models_from_trainers=True,   # 参数由 Trainer 初始化同步
        buffer_cpu_offload=True,
        kl_coef=0.1,
    )
    for i in range(num_makers)
]

detached_trainer_name_list 中的名字就是该 Maker 要把经验发往的目标 Trainer;文档提示 Trainer 的名字与 Maker 命名方式一一对应(trainer_{i} / maker_{i}),通过 .options(name=...) 设定。

3.3 Setup Trainers(配置训练端)

第一步:定义 Trainer 的环境变量(与 Maker 结构相同,但端口与 world size 独立):

env_info_trainers = [{
    'local_rank': '0',
    'rank': str(rank),
    'world_size': str(num_trainers),
    'master_port': trainer_port,
    'master_addr': master_addr
} for rank in range(num_trainers)]

第二步:定义 Trainer 侧模型(只需 actor 与 critic):

def trainer_model_fn():
    actor = get_actor_from_args(model, pretrained)
    critic = get_critic_from_args(model, pretrained)
    return actor, critic

第三步:创建 Trainer 命名 Actor(注意 model_fn 应传入工厂函数本身,而非其返回值):

trainer_refs = [
    DetachedPPOTrainer.options(
        name=f"trainer_{i}",
        num_gpus=1,
        max_concurrency=2
    ).remote(
        experience_maker_holder_name_list=[f"maker_{x}" for x in target_makers(i, num_trainers, num_makers)],
        strategy_fn=lambda: get_strategy_from_args("colossalai_zero2"),
        model_fn=trainer_model_fn,
        env_info=env_info_trainers[i],
        train_batch_size=8,
        buffer_limit=0,
        eps_clip=0.2,
        value_clip=0.4,
    )
    for i in range(num_trainers)
]

experience_maker_holder_name_list 指明该 Trainer 训练后要回传参数给哪些 Maker。两个 name list 组合起来就定义了完整的 Maker↔Trainer 传输图(例如 2:1 时每个 Trainer 同时收到两个 Maker 的经验,并回传给两个 Maker)。

策略工厂 get_strategy_from_argsutils.py#L70-L85)支持五种策略,覆盖文档所说的"推理与训练可各取所需":

策略名 实际构造 定位
ddp DDPStrategy() 数据并行基线
colossalai_zero2 LowLevelZeroStrategy(stage=2, placement_policy="cuda") Zero-2,GPU 侧放置
colossalai_gemini GeminiStrategy(placement_policy="static", initial_scale=2**5) Gemini 混合内存调度
colossalai_gemini_cpu Gemini + offload_optim_frac=1.0, offload_param_frac=1.0 参数/优化器全量 CPU offload
colossalai_zero2_cpu LowLevelZeroStrategy(stage=2, placement_policy="cpu") Zero-2,CPU 侧放置

从源码结构看,一个常见组合是:Maker 用 colossalai_gemini(推理吞吐优先),Trainer 用 colossalai_zero2colossalai_gemini_cpu(省显存、可跑更大 batch),这正是文档"Flexible Structure"一节的落地方式。

3.4 Launch Jobs(启动训练)

定义数据加载器工厂(在 Maker 端执行,因此闭包内不要提前构造):

def data_loader_fn():
    return torch.utils.data.DataLoader(dataset, batch_size=experience_batch_size,
                                        collate_fn=collate_fn)

启动 Maker 工作循环与 Trainer 训练,然后统一等待:

wait_tasks = []
for ref in experience_holder_refs:
    wait_tasks.append(ref.workingloop.remote(data_loader_fn, num_steps=experience_steps))
for ref in trainer_refs:
    wait_tasks.append(ref.fit.remote(total_steps, update_steps, train_epochs))
ray.get(wait_tasks)

对照源码可以确认这两个入口的语义:

  • workingloop(dataloader_fn, num_epochs=1, num_steps=0)experience_maker_holder.py#L141-L168):num_steps > 0 时按步数迭代(迭代器耗尽后自动重开),否则按 epoch 遍历;循环开始前会 _get_ready() 自旋等待初始化完成(当 sync_models_from_trainers=True 时,需等首个 Trainer 完成 fully_update 参数同步)。
  • fit(total_steps, update_steps, train_epochs=1)detached_trainer_base.py#L105-L114):按 update_steps 个训练步为一轮,每轮结束回传参数给目标 Maker。

由于两侧都是命名 Actor 且全部通过异步 RPC 通信,ray.get(wait_tasks) 返回即代表整个 PPO 流水线(经验生产、缓冲、训练、参数回传)全部结束。

4. Flexible Structure:弹性部署形态

文档用四幅示意图(托管于项目外部资产库,此处以文字描述)给出了四种典型拓扑,均可通过第 3 节的 name list 与策略参数组合实现:

  1. 2 Makers : 1 Trainer:两个 Maker 并发推理产经验,轮询投递给同一个 Trainer;Trainer 回传参数时同时推送给两个 Maker。推理吞吐与训练吞吐解耦,Maker 数可以独立扩容。
  2. 2 Makers : 2 Trainers:Maker 与 Trainer 一一对应(或交叉连接),传输图更均衡,get_receivers_per_sender 可按取模规则自动生成连接关系。
  3. Maker 推理量化:Maker 侧加载 INT8/量化推理模型加速采样,Trainer 侧保持 FP16 训练;由于参数回传走 state_dict 分片协议,Maker 侧只需在 update_experience_maker 中反量化后再装载即可。从源码结构看,该能力主要依赖 Maker 的 model_fn 自定义加载逻辑,与传输协议正交。
  4. Tensor Parallel:Maker 或 Trainer 内部再叠加 TP(如每 2 GPU 一组做张量并行),外部仍是 Maker↔Trainer 的解耦图。文档 TODO 中"Support TP & PP"说明原生 TP/PP 策略当时仍在规划中,实践中需自行通过 strategy_fn 组合相应并行策略。

5. 关键注意点与文档时效性

  • 文档时效README 原文 首行明确警告"此内容在 ColossalChat 大更新后可能过时"。因此示例路径 applications/Chat/examples/ray 应理解为现在的 examples/community/ray;示例脚本使用旧式接口(ExperienceMakermake_experience 手动组装 Experience),生产化代码建议直接基于 coati/ray 模块的现成封装。
  • LoRA 状态:文档 TODO 写着"Support LoRA"未完成,但当前源码中 update_lora_weightsLoRAConstructorfilter_state_dict_lorareconstruct_increase 等链路均已存在于 Maker 与 Trainer 两侧,可推断 LoRA 增量同步已具备雏形实现,只是尚未作为完整文档化能力提供。
  • 参数同步正确性依赖 chunk 协议chunk_start/chunk_end 的加锁顺序保证 Maker 在参数分片写入期间暂停生成;若自定义 Maker 逻辑时绕过该协议(例如自行 load 新参数),需要自行处理模型访问锁(_model_visit_lock)。
  • 经验分片:Maker 发送前用 split_experience_batch 把批拆成单条 item 再轮询投递,因此 Trainer buffer 中的粒度是单条经验;buffer_extendbuffer_sample 分属不同并发组,写入与采样互不阻塞。

6. 小结

ColossalChat 的 Ray 模块用 ExperienceMakerHolderDetachedPPOTrainer 两个 Ray Actor、以及命名 Actor + 并发组的组合,把 PPO 的推理侧与训练侧解耦成可独立伸缩的异步流水线:Maker 持续采样并经轮询投递经验,Trainer 消费 buffer 训练后按分片协议回传参数,两侧可分别套用 Zero/Gemini、CPU offload、量化推理等差异化策略。相关实现集中在 coati/ray 模块,配套示例位于 examples/community/ray,可沿第 3 节的步骤逐段复现。

登录后查看全文
热门项目推荐
相关项目推荐