ColossalAI ColossalChat 基于 Ray 的 Detached PPO 分布式训练:Maker 与 Trainer 异步解耦实战指南
本文围绕 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),这同样有利于横向扩展。
文档明确指出,DetachedPPOTrainer 与 ExperienceMakerHolder 是 Ray 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_step(experience_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_maker(experience_maker_holder.py#L170-L229)采用分块(chunked)协议:
- Trainer 先发
chunk_start=True标记,Maker 侧acquire()模型访问锁——此刻推理暂停; - Trainer 循环发送
new_actor_state_dict/new_critic_state_dict的各个分片,Maker 逐片load_state_dict(..., strict=False); - 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_makers,compute 承载 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 增量 |
注意一个实现细节:当策略是 LowLevelZeroStrategy 或 GeminiStrategy 时,Trainer 自动选用 ColossalAI 的 HybridAdam 优化器,否则用 torch.optim.Adam(detached_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_sender(utils.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 示例说明,标准流程是:
- 按 Ray 官方指南搭建支持 GPU 的 Ray 集群,记录 API server 端点(形如
http://your.head.node.address:8265); - 克隆本仓库;
- 通过 Job Submission Client 提交任务:
python applications/ColossalChat/examples/community/ray/ray_job_script.py http://your.head.node.address:8265
- 打开 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_info(utils.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_args(utils.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_zero2 或 colossalai_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 与策略参数组合实现:
- 2 Makers : 1 Trainer:两个 Maker 并发推理产经验,轮询投递给同一个 Trainer;Trainer 回传参数时同时推送给两个 Maker。推理吞吐与训练吞吐解耦,Maker 数可以独立扩容。
- 2 Makers : 2 Trainers:Maker 与 Trainer 一一对应(或交叉连接),传输图更均衡,
get_receivers_per_sender可按取模规则自动生成连接关系。 - Maker 推理量化:Maker 侧加载 INT8/量化推理模型加速采样,Trainer 侧保持 FP16 训练;由于参数回传走 state_dict 分片协议,Maker 侧只需在
update_experience_maker中反量化后再装载即可。从源码结构看,该能力主要依赖 Maker 的model_fn自定义加载逻辑,与传输协议正交。 - 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;示例脚本使用旧式接口(ExperienceMaker、make_experience手动组装 Experience),生产化代码建议直接基于coati/ray模块的现成封装。 - LoRA 状态:文档 TODO 写着"Support LoRA"未完成,但当前源码中
update_lora_weights、LoRAConstructor、filter_state_dict_lora、reconstruct_increase等链路均已存在于 Maker 与 Trainer 两侧,可推断 LoRA 增量同步已具备雏形实现,只是尚未作为完整文档化能力提供。 - 参数同步正确性依赖 chunk 协议:
chunk_start/chunk_end的加锁顺序保证 Maker 在参数分片写入期间暂停生成;若自定义 Maker 逻辑时绕过该协议(例如自行 load 新参数),需要自行处理模型访问锁(_model_visit_lock)。 - 经验分片:Maker 发送前用
split_experience_batch把批拆成单条 item 再轮询投递,因此 Trainer buffer 中的粒度是单条经验;buffer_extend与buffer_sample分属不同并发组,写入与采样互不阻塞。
6. 小结
ColossalChat 的 Ray 模块用 ExperienceMakerHolder 与 DetachedPPOTrainer 两个 Ray Actor、以及命名 Actor + 并发组的组合,把 PPO 的推理侧与训练侧解耦成可独立伸缩的异步流水线:Maker 持续采样并经轮询投递经验,Trainer 消费 buffer 训练后按分片协议回传参数,两侧可分别套用 Zero/Gemini、CPU offload、量化推理等差异化策略。相关实现集中在 coati/ray 模块,配套示例位于 examples/community/ray,可沿第 3 节的步骤逐段复现。
atomcodeClaude Code 的开源替代方案。连接任意大模型,编辑代码,运行命令,自动验证 — 全自动执行。用 Rust 构建,极致性能。 | An open-source alternative to Claude Code. Connect any LLM, edit code, run commands, and verify changes — autonomously. Built in Rust for speed. Get StartedRust0623
Hy4-previewHy4 preview 是由腾讯混元团队研发的新一代混合专家(MoE)旗舰模型。模型总参数量 770B,每个 token 激活 49B,主干共包含78层,第一层采用标准 FFN,其余 77 层均为 MoE 结构,每层包含 256 个路由专家与 1 个共享专家,每个 token 激活 top-8 路由专家及共享专家。主干之外原生内置 1 层 MTP(总参数量 10B,激活 0.7B)以支持投机解码。Python00
GLM-5.3GLM-5.3 与 GLM-5.2 使用相同的基座模型——所有提升均来自后训练。与 GLM-5.2 相比,它在复杂编程和长程任务上的表现显著提升。Jinja00
GLM-5.3-FlashGLM-5.3-Flash (320B-A18B),是GLM-5系列的首个原生多模态模型。320B总参数,能力超过GLM-5.2Jinja00
Spark-X2.5-4BSpark-X2.5-4B 旨在让强大的 AI 更实用、更高效、更易获得。在广泛日常任务中表现强劲,涵盖对话、写作、翻译、推理、编码、工具调用以及智能体工作流,并在同等规模的开源模型中取得领先成绩。Spark-X2.5 将面向效率的架构与最高 1M tokens 的原生上下文窗口相结合,并支持 200 多种语言。Python00
Spark-X2.5-1.7BSpark-X2.5-1.7B 旨在让强大的 AI 更加实用、高效且易于获取。这些模型在广泛的日常任务中表现出色,涵盖对话、写作、翻译、推理、编程、工具调用和智能体工作流,并在同等规模的开源模型中取得领先结果。Spark-X2.5 将面向效率的架构与最高 1M tokens 的原生上下文窗口相结合,并支持 200 多种语言。Python00