Agent Lightning 异步训练详解:Collocated Asynchronous Rollout 的开启、原理与调优
【免费下载链接】agent-lightningThe absolute trainer to light up AI agents.项目地址: https://gitcode.com/GitHub_Trending/ag/agent-lightning
长时程 Agent 的 rollout 时长差异很大,同步训练中一个慢 rollout 就会拖住整个更新步。Agent Lightning v1.0 提供的collocated 异步训练让 rollout 生成与模型更新共享同一 GPU 池,未完成的任务组自动结转(carry over)到后续步骤。读完本文,你将掌握如何用两行配置开启异步收集、理解“暂停/排空 Gateway”的底层机制、读懂 W&B 上报的六项异步指标,并配置 rollout correction 修正策略失配(staleness)。
为什么需要异步训练
在同步模式下,一次参数更新必须等待该步所有 rollout 组全部完成,更新步的耗时等于最慢的那个组。对于需要多轮工具调用、检索或沙箱执行的任务,rollout 时长方差极大,GPU 会大量空转等待。Agent Lightning 的解法是collocated(共置)异步:推理与训练不再分置不同设备,而是共享 GPU 池,训练侧只要凑够一批完成的组就立即更新,其余组继续跑,同时把新的 prompt 组补进活跃队列。
开启异步训练
通过agentlightning.async_rollout.enabled打开异步收集,并且必须显式设置async_train_batch_size:
agentlightning: async_rollout: enabled: true async_train_batch_size: 64async_train_batch_size是保持活跃(active)的 prompt 组数量,它必须严格大于data.train_batch_size——后者是单次更新消费的已完成组数量。完整的配置示例:
data: train_batch_size: 32 agentlightning: async_rollout: enabled: true async_train_batch_size: 64这个约束在 Trainer 初始化时会被强制校验:trainer.py 中,若async_train_batch_size <= train_batch_size会直接抛出ValueError。仓库自带的默认配置见 config.yaml,其中enabled: false、async_train_batch_size: null,即默认走同步模式。
一个实用的起步值:
$$B_{async} = 2 B_{train}$$
即async_train_batch_size取train_batch_size的两倍。仓库中的官方示例也是这么做的:train_calc_agent.py 在--async模式下自动将async_train_batch_size设为train_batch_size * 2;train_gsm8k_agent.py、train_sw_agent.py、train_smith_agent.py 等示例同样在配置中预置了该项。
调优方向(沿用官方文档的建议):
- 调大
async_train_batch_size:当 rollout 时长方差显著、且承载 Agent 的资源(本地进程或 Kubernetes Job)容量充足时,更大的活跃窗口能更充分地掩盖慢任务。 - 调小:当活跃的 Agent 进程 / K8s Job 消耗过多 CPU 或内存资源时,应收缩窗口以控制并发压力。
工作机制:活跃窗口、结转与 Gateway 排空
官方文档将异步收集描述为六个步骤,下面结合源码逐步印证。
Trainer 维持至多
async_train_batch_size个活跃 prompt 组。见 trainer.py 的_next_train_batch_dict_for_rollout:每步先统计上一轮结转的组数n_carry_over,再按n_new = async_train_batch_size - n_carry_over从 dataloader 取新组,从而把活跃窗口恒定在B_async。Controller 在本地进程或 Kubernetes Job 中启动 Agent 执行。rollout 通过
AglAsyncRolloutManager(agl_rollout_manager.py)批量提交给 Controller,由后者按local.agent_class或k8s.job_template_path配置的模板拉起。凑满
data.train_batch_size个完成的组即用于更新,不等所有活跃组。核心是AglAsyncRolloutManager.enqueue_and_wait_until_group_completed(agl_rollout_manager.py):它把结转 rollouts 与新 rollouts 合并为active_rollouts,按data_id分组后轮询状态,while len(completed_group_keys) < target_finished_group_num循环直到完成组数达到train_batch_size(Trainer 侧在 trainer.py 以target_finished_group_num=self.config.data.train_batch_size传入)。未完成的组保持活跃并结转到下一步。循环结束后,所有未完成的组打包为
new_carry_over_rollouts返回,Trainer 将其存入self._carry_over_rollouts(trainer.py),供下一轮补位。更新权重前,Gateway 暂停新请求并等待在途请求排空。完成组收集后,
_rollout会调用_pause_and_drain_gateway(trainer.py):先POST /proxy/pause(携带原因,超时放宽到 300 秒),拿到暂停瞬间的inflight计数,然后每 0.25 秒轮询GET /proxy/state,直到inflight <= 0,并记录排空耗时。服务端实现在 routes/proxy.py 的/pause、/resume、/state三个管理端点,状态由 proxy.py 中的ProxyPauseState(paused、inflight、reason等字段,带asyncio.Lock)维护。共享 GPU 执行模型更新,随后推理恢复,进入下一个 rollout 阶段。下一次
_rollout开始时,Trainer 会先_resume_gateway(trainer.py,幂等可重试)恢复转发。
组完整性:GRPO/RLOO 的统计前提
每个 prompt 组必须保持完整:例如actor_rollout_ref.rollout.n为 4 时,该 prompt 的 4 个兄弟 rollout必须全部结束,这个组才能被优化器消费。源码中这一点有硬断言——AglAsyncRolloutManager按data_id分组后检查assert len(group) == self._train_rollout_n(agl_rollout_manager.py),并只在组内所有 rollout 都到达终态时才计入completed_group_keys(L503-L516)。这样做保证了组内兄弟样本来自同一 prompt,GRPO/RLOO 等依赖组内对比的优势估计不被破坏。
Agent 侧要求:使用带重试的客户端
当 Agent 的请求恰好落在 Gateway 暂停窗口时,forward_request 会立即返回一个可重试响应:HTTP 429,附Retry-After头(默认 5 秒,见ProxyPauseState.retry_after_seconds)与X-Agl-Paused: true标记,body 为{"error": "gateway paused", "reason": ...}。因此官方建议 Agent 使用带重试的 OpenAI 或 HTTP 客户端:暂停期间到达的请求收到 429 后自动退避,推理恢复后即可继续,整个 Agent 流程无需感知训练更新。
监控:W&B 异步指标
Trainer 每步向 W&B 上报一组training/async/前缀的指标,其计算逻辑见_compute_async_rollout_metrics,加上暂停/排空阶段的两个指标(trainer.py):
| 指标 | 含义 |
|---|---|
training/async/n_prev_carry_over_rollouts | 从上一步继承过来的 rollout 数 |
training/async/n_completed_rollouts | 本步被消费(完成)的 rollout 数 |
training/async/n_new_carry_over_rollouts | 结转给下一步的未完成 rollout 数 |
training/async/new_carry_over_age_max_steps | 当前结转 rollout 中“最老”的已跨过的优化器步数 |
training/async/proxy_inflight_at_pause | Gateway 暂停时刻仍在执行的请求数 |
training/async/proxy_drain_seconds | 等待在途请求全部结束所花的时间 |
判读建议:n_new_carry_over_rollouts长期逼近窗口上限、new_carry_over_age_max_steps持续增大,说明慢任务积压,可考虑调大async_train_batch_size或收紧 rollout 超时(rollout_timeout_seconds,默认 1800 秒,见 config.yaml);proxy_drain_seconds反映每次更新前排空在途请求的代价,通常应保持在一个较小量级。此外 Trainer 还会上报timing/rollout_*系列排队/运行耗时聚合(trainer.py),可用于观察 K8s 批量拉起 Pod 造成的队列等待。
处理 Staleness:rollout correction
异步 rollouts 由较旧版本的模型生成,可能在被用于训练前就已“过期”(stale),带来策略失配。对此,官方方案是启用verl的rollout correction,并推荐token 级重要性采样(TIS)、裁剪阈值设为2:
algorithm: rollout_correction: rollout_is: token rollout_is_threshold: 2该配置直接并入 verl 的 PPO trainer 配置,与agentlightning.async_rollout相互独立:前者解决“用旧策略采的样本如何纠正”,后者解决“如何不停机地持续采样”。
适用前提与限制小结
- 适用 Agent Lightningv1.0(当前仓库主干);v1.0 之前的旧版架构(v0.x 分支)不适用本文配置项。
- 必须是共置部署:rollout 生成与模型更新共享同一 GPU 池,暂停/排空机制依赖 API Gateway 的
/pause与/state端点。 - 开启后必须同时设置
async_train_batch_size且严格大于train_batch_size,否则 Trainer 初始化即报错。 - Agent 侧应使用可处理 429 +
Retry-After的重试客户端,否则会在暂停窗口内中断。 - 建议同时评估 staleness 处理:rollout 越多、跨步结转越久,越应启用
rollout_correction。
配套阅读:Trainer Configuration(verl 集成与 trace 聚合)、API Gateway Configuration(Gateway 与模型代理配置)、Controller Configuration(本地与 K8s 执行器);行为与约束的交叉引用见 4-trainer-configuration.md 末尾对本文档的回链。
【免费下载链接】agent-lightningThe absolute trainer to light up AI agents.项目地址: https://gitcode.com/GitHub_Trending/ag/agent-lightning
创作声明:本文部分内容由AI辅助生成(AIGC),仅供参考