文档教程人工智能大模型RLHF【免费下载链接】Awesome-ML-SYS-TutorialMy learning notes for ML SYS.项目地址https://gitcode.com/gh_mirrors/aw/Awesome-ML-SYS-Tutorial点击查看免费下载导读本指南以 over_sample.md 为核心系统讲解 verl 与 SGLang 多轮multi-turnrollout 中为解决长尾long-tail问题而引入的Over-Sample过采样特性它如何在每个 rollout worker 上超额发起请求、在完成目标数量的请求后立即 abort 剩余请求并替换为 padding 数据。读完本文你将掌握从 Docker 环境搭建、over_sample分支源码安装、GSM8K / DAPO 训练复现的完整实操流程并深入理解monitor_and_cancel、process_request_with_monitoring、run_with_cancellation三个核心函数的协作机制以及 padding 对 loss、reward 指标和 GRPO group 带来的影响与待办事项。一、背景为什么多轮 rollout 需要 Over-Sample1.1 长尾问题是异步多轮 rollouts 的核心瓶颈在多轮 RLHF / Agentic RL 训练中同一个 batch 内的请求会经历多轮对话与工具调用每个请求的轮数、生成长度差异极大。来自仓库 profiling 分析 profile.md 的实测结论指出长尾效应非常显著在一个 step 的 rollout 过程中80% 的请求会在前 40%50% 的时间内完成其余请求则拖着长长的尾巴同步 rollout 中batch 中最慢的样本会拖慢整个流水线参考 verl-multiturn-rollout-Release.md 中对异步 rollout 动机的说明长尾请求往往是轮数更多、难度更大的样本其响应长度会急剧增长而对应的 reward 通常更低。正是因为多轮对话 工具调用会让请求的耗时方差急剧放大业界才普遍采用over sample过采样策略来对冲长尾多发起一批请求只要凑够足够数量的真实完成样本就放弃剩余未完成的请求从而把「等待最慢请求」的时间压缩为「等待目标数量请求」的时间。1.2 Over-Sample 与 partial rollout 的区别profile.md 与 profile_en.md 中明确对比了两类方案partial rollout把过采样但未完成的请求保存下来下一个 step 在上一 step 的基础上继续 rollout属于细粒度的异步化与续跑方案over sample本文方案比 partial rollout更粗暴未完成的请求直接被丢弃不续跑、不保留轨迹只替换为 padding 数据。此外profile_en.md 还记录了一条更朴素的替代路线——直接调低max response length例如从 1000 降到 600也能显著压缩每 step 耗时且不影响收敛这从侧面印证了长尾请求对训练时间的高贡献与对 reward 的低贡献。Over-Sample 则是在不牺牲长响应能力的前提下用「超额发起 提前截止」的思路换取训练吞吐。二、快速复现从 Docker 到 8 卡训练以下复现步骤完整继承自 over_sample.md并与同目录 latest_sglang.md 的环境搭建经验相互印证。2.1 创建 Docker 容器使用前需要配置好WANDB_API_KEY获取方式参考 WandB 官方文档「Where can I find the API token」。如果你的系统尚未配置HF_TOKEN与WANDB_API_KEY请先导出这两个环境变量再启动容器docker run -it --name h100_verl_{your_name} --gpus all \ --shm-size 32g \ -v {your_cache_path}:/root/.cache \ --env HF_TOKEN$HF_TOKEN \ --env WANDB_API_KEY$WANDB_API_KEY \ --ipchost \ lmsysorg/sglang:latest \ /bin/bash进入容器后可以确认环境变量是否被正确映射echo $HF_TOKEN echo $WANDB_API_KEY以后每次从容器exit出来后用下面这条命令即可重启容器容器名、挂载、环境变量均保持不变docker start -i h100_verl_{your_name}2.2 配置 Python 环境并基于源码安装 verl-sglangmkdir -p /tmp chmod 1777 /tmp apt update apt install -y python3.10 python3.10-venv python3 -m ensurepip --upgrade python3 -m venv ~/.python/verl-sglang source ~/.python/verl-sglang/bin/activate python3 -m pip install --upgrade pip python3 -m pip install --upgrade uv接着克隆带over_sample功能的 verl 分支并安装cd ~ git clone -b over_sample https://github.com/zhaochenyang20/verl.git cd verl python -m uv pip install wheel setuptools python3 -m uv pip install -e .[sglang] --prereleaseallow python3 -m uv pip install -r ./requirements.txt --no-build-isolation python3 -m uv pip install torch_memory_saver说明--prereleaseallow允许安装预发布版本依赖--no-build-isolation用于在已有构建环境内安装依赖避免重复编译torch_memory_saver用于 rollout engine 与训练 engine 共享 GPU 时的显存释放/恢复对应后文release_memory_occupation/resume_memory_occupation机制。2.3 测试 GSM8K 多轮训练cd ~/verl export CUDA_VISIBLE_DEVICES0,1,2,3,4,5,6,7 # 拉取并预处理 gsm8k 数据集生成带工具调用要求的 prompt 与 ground truth python examples/data_preprocess/gsm8k_multiturn_w_tool.py # 启动 8 卡训练 bash examples/sglang_multiturn/run_qwen2.5-3b_gsm8k_multiturn.shgsm8k_multiturn_w_tool.py的核心逻辑参见 readme.md包括加载openai/gsm8k原始数据、为每条样本生成带有工具调用要求的 prompt如要求模型调用calc_gsm8k_reward工具、把 ground truth 写入extra_info字段、最后存储为train.parquet/test.parquet。2.4 测试 DAPO 多轮训练cd ~/verl export CUDA_VISIBLE_DEVICES0,1,2,3,4,5,6,7 bash examples/sglang_multiturn/run_qwen3_4b_dapo_multiturn.sh需要特别注意的是DAPO 训练依赖 sandbox-fusion 工具服务器。参考 latest_sglang.md需要在宿主机上单独启动工具服务器# 启动 sandbox fusiondapo tool call requirement docker run -it -p 8080:8080 volcengine/sandbox-fusion:server-20250609同时为了让训练容器能访问宿主机 sandbox-fusion 的 8080 端口第一步创建容器时需要额外加上--networkhostdocker run -it --name h100_verl_{your_name} --gpus all \ --shm-size 32g \ -v {your_cache_path}:/root/.cache \ --env HF_TOKEN$HF_TOKEN \ --env WANDB_API_KEY$WANDB_API_KEY \ --networkhost \ --ipchost \ lmsysorg/sglang:latest \ /bin/bash环境提示来自 latest_sglang.md 的踩坑记录如果启动后遇到形如ValueError: Feature type List not found的报错请向上翻看完整调用栈——真正的根因往往是Python 环境错位主进程使用了虚拟环境/root/.python/verl-sglang/lib/python3.10/site-packages/而 Ray worker 进程却用了系统 Python/usr/local/lib/python3.10/dist-packages/。作者最终的建议是干脆不要使用虚拟环境逆天 ray...。三、设计思路三个函数协同完成「超额采样 提前截止」Over-Sample 特性基于 commitb979a73e358313afafab5db512cd5ae0009ccac0实现。整体设计用一个句子可以概括同时启动所有请求的异步 rollout并持续监控「真实完成数量」一旦达到目标完成数target_completion立即取消剩余任务、向 SGLang engine 发送 abort 信号把未完成请求统一替换为 padding 数据。整个机制由三个函数协作完成函数职责process_request_with_monitoring处理单个请求真实完成则计数并返回真实结果目标已达成后完成的请求返回 paddingmonitor_and_cancel持续监控完成数量达到目标后取消剩余任务并向 engine 发送 abort 信号run_with_cancellation同时启动上述两者收集所有结果把异常统一转换为 padding三者共享两个全局变量completed_count累计真实完成数与completion_lock保证计数原子性的读写锁。3.1 前置数据total_requests本 worker 本次 step 发起的总请求数超额采样的总量target_completion目标完成数小于total_requests二者之差即「允许丢弃的尾部请求预算」每个请求对应一个asyncio.Task其内部阻塞等待_async_rollout_a_request完成。四、源码级解析三个核心函数4.1process_request_with_monitoring每个请求独立地「真实完成 or padding」async def process_request_with_monitoring(req): nonlocal completed_count try: result await self._async_rollout_a_request(req, do_sample, is_validate, **kwargs) async with completion_lock: if completed_count target_completion: completed_count 1 print(f✅ Request {req.request_id} completed ({completed_count}/{total_requests})) return result # 返回真实结果 else: # 超过目标返回padding logger.info(fRequest {req.request_id} finished after target met, creating padding) return self._create_padding_request(req)执行语义拆解每个请求独立启动每个 request 在自己的process_request_with_monitoring任务中通过await阻塞式执行_async_rollout_a_request因此不同请求天然异步并发较早完成的请求其result是真实 rollout 结果并递增completed_count。由于completed_count是全局共享变量必须通过completion_lock保证计数操作的原子性避免并发读写冲突较晚完成的请求当monitor_and_cancel检测到completed_count已达到target_completion后这些任务会被取消并向 SGLang engine 发送abort_requests若某个请求在目标达成之后、被取消之前刚好完成则走else分支返回_create_padding_request(req)的 padding 结果。4.2monitor_and_cancel监控 取消 向 engine 发 abortasync def monitor_and_cancel(): nonlocal completed_count while completed_count target_completion: await asyncio.sleep(0.1) # 每0.1秒检查一次 print(f Target reached: {completed_count}/{total_requests} completed!) print( Cancelling remaining requests and sending abort to engine...) # 取消剩余的任务 cancelled_count 0 for task in all_tasks: if not task.done(): task.cancel() cancelled_count 1 # 向engine发送abort信号 try: abort_result await self._engine.abort_request(abort_allTrue) print(f✅ Abort signal sent to engine: {abort_result}) except Exception as e: print(f❌ Failed to send abort signal to engine: {e})要点轮询粒度每 0.1 秒检查一次completed_count一旦达到target_completion立即行动取消任务遍历all_tasks对未完成not task.done()的任务调用task.cancel()触发asyncio.CancelledErrorengine 级 abort调用self._engine.abort_request(abort_allTrue)让 SGLang engine 层面也释放/停止这些请求的生成避免被取消的请求继续占用显存与算力。从 profile.md 的 profiling 视角看被 abort 的请求会记录aborted_request_with_cancelled_error事件且被 abort 的请求耗时高度一致详见第七节符合「到达目标时间点后统一截止」的预期。4.3run_with_cancellation并发编排与结果兜底async def run_with_cancellation(): nonlocal all_tasks # 创建所有任务 all_tasks [asyncio.create_task(process_request_with_monitoring(req)) for req in req_list] # 启动监控任务 monitor_task asyncio.create_task(monitor_and_cancel()) try: # 等待所有任务完成包括被取消的 results await asyncio.gather(*all_tasks, return_exceptionsTrue) # 处理结果将异常转换为padding output_req_list [] for i, result in enumerate(results): if isinstance(result, Exception): # 异常转换为padding logger.warning(fTask {i} resulted in exception: {result}) output_req_list.append(self._create_padding_request(req_list[i])) else: output_req_list.append(result) return output_req_list finally: # 清理监控任务 monitor_task.cancel() try: await monitor_task except asyncio.CancelledError: pass理解要点all_tasks与读写锁completion_lock是三个函数的全局变量这里同时启动所有 reqs 的process_request_with_monitoring并额外创建一个monitor_task来监视完成进度关键洞察虽然每个请求的_async_rollout_a_request未必能完成但上层的process_request_with_monitoring一定会结束要么返回真实结果 / padding要么以CancelledError等异常结束。因此results await asyncio.gather(*all_tasks, return_exceptionsTrue)一定会返回逐个处理results时存在三种情况COMPLETED真实结果、Exception被取消或出错、PADDING目标达成后才完成的。将Exception转换为PADDING后返回output_req_list最后在finally中取消并回收monitor_task保证不留悬挂任务。4.4 调用入口在实际的 rollout worker 中参见 profile.md 记录的代码结构异步 rollout 逻辑位于上述三个函数定义之后通过事件循环驱动# run async tasks self.log_path os.path.join(self.log_dir, fstep_{self.step}, fworker_{self._rank}.jsonl) torch.cuda.synchronize() async_rollout_with_monitoring_start_time time.time() loop asyncio.get_event_loop() output_req_list loop.run_until_complete(run_with_cancellation()) torch.cuda.synchronize() async_rollout_with_monitoring_end_time time.time()因此async_rollout_with_monitoring_start_time与async_rollout_with_monitoring_end_time之间的时长就是该 worker 上所有请求含被 abort 的完成本轮 rollout 的总耗时可直接用于衡量 Over-Sample 对 rollout 耗时的压缩效果。五、engine 侧AsyncEngine.abort_request与同步/异步的边界文档指出monitor_and_cancel中发送的 engine abort 信号其实际实现在sglang_rollout.py的AsyncEngine类中async def abort_request(self, rid: str , abort_all: bool False): Abort a specific request or all requests. Args: rid: The request ID to abort. If empty and abort_all is False, no action is taken. abort_all: If True, abort all running requests regardless of rid. try: result self.tokenizer_manager.abort_request(ridrid, abort_allabort_all) print(f Abort result: {result}) return result if result is not None else {status: aborted} except Exception as e: logger.error(fFailed to abort requests: {e}) raise这里有几点值得深入玩味verl 的AsyncEngine继承并重写了 SGLang Engine 的许多方法比如update_weights_from_tensor和resume_memory_occupation。从架构上看SGLang Engine 不实现这些方法也不影响 verl但会影响其他框架。最初作者以为必须先让 SGLang Engine 原生支持abort_request因为起初只有 server 端有而 engine 没有但既然AsyncEngine已经重写了abort_requestSGLang Engine 无需实现该功能也无需为此发版——毕竟在 verl 上更新 SGLang 版本代价很高。同步/异步的边界问题与update_weights_from_tensor不同abort_request内部不能用await去调用self.tokenizer_manager.abort_request必须直接调用。这取决于 SGLangtokenizer_manager内部的实现如果某函数在 tokenizer_manager 中是异步实现的外部调用才能使用await语法。令人费解的是resume_memory_occupation与abort_request同在 tokenizer_manager 中前者是异步的、后者却是同步的。「异步函数里只写一行 await」的意义作者提出疑问——在一个异步函数中仅await另一个异步函数本质上是在等待内层异步函数执行完成那么外层函数写成同步是否也行答案是否定的因为外层调用方如fsdp_sglang.py中的release_memory也需要await外层必须是异步函数才能被await。来看sglang_rollout.py中resume_memory_occupation的实现async def resume_memory_occupation(self, tags: Optional[list[str]] None): Resume GPU occupation. # because __init__ is a sync method, it can not call the async release_memory_occupation # have to move release_memory_occupation from __init__ to here # For multi-stage awake, we run release weight and kv_cache when we resume weights for the first time. if self._need_reload: await self.release_memory_occupation() self._need_reload False if tags is None: obj ResumeMemoryOccupationReqInput() else: obj ResumeMemoryOccupationReqInput(tagstags) return await self.tokenizer_manager.resume_memory_occupation(obj, None)注释点明了一个关键约束__init__是同步方法无法调用异步的release_memory_occupation所以必须把该调用从__init__挪到这里对多阶段唤醒multi-stage wake-up场景第一次恢复权重时需要同时释放 weight 与 kv_cache通过self._need_reload标志控制而fsdp_sglang.py中release_memory的调用链如下展示了外层如何awaitasync def release_memory(self): if self.device_mesh[infer_tp].get_local_rank() 0 and self.rollout_config.free_cache_engine: if self.multi_stage_wake_up: await self.inference_engine.release_memory_occupation(tags[kv_cache, weights]) else: await self.inference_engine.release_memory_occupation() log_gpu_memory_usage(After release memory occupation in sharding manager, loggerlogger)从中可以推断异步函数的「可等待性」由最内层的 tokenizer_manager 实现决定——若内层是同步函数如abort_request外层直接调用即可若内层是异步函数如resume_memory_occupation外层就必须通过await等待其完成并沿着调用链逐层保持异步签名。六、Padding 的正确性争议与待办清单作者在整体读完实现后给出了克制而务实的评价「设计的还算清晰实现可能未必好还要大改」。并明确列出了必须检查的地方6.1 直接打成 padding 是否真正「丢弃」了请求理想设计中被 abort 的请求应被视为不存在GRPO 的 group size 相应减小还能省下训练时间。为此需要仔细核查除将response_loss_mask设为 0 之外是否还有其他需要修改的地方作者最初修改了agg_loss函数但咨询后认为可能并不需要需额外确认reward 函数是否需要更改对于 FSDP对应_expand_to_token_level一个待验证的风险如果实现了「完美丢弃」不同 GRPO group 的请求数量将不一致理论上会影响 GRPO group 的方差可能导致训练更不稳定——这一影响目前无法定量刻画。6.2 另一种设计保留轨迹、mask 掉 loss直接打成 padding 相当于把 partial rollout 得到的 trajs 直接丢了。另一种更温和的设计是保留这些 trajs但 loss mask 设为 0reward 也设为 0。文档中引用了龙老师的观点认为这种「保留但不参与学习」的方案可能更好。该方案的具体效果对 loss、对 reward 均值、对 GRPO 方差同样留待验证。6.3 reward 指标虚高问题与compute_data_metrics的改动整个 feature 需要修改sglang_rollout.py与metric_utils.py的compute_data_metrics函数。后者非常 tricky当前实现把被 abort 的请求在每次 reward 均值计算时排除在分母之外。由此带来一系列问题validation 阶段无 abort理论上 validation step 不会产生被 abort 的请求因此不会受到aborted_mask (response_length 0).bool()的影响但仍需实验验证reward 虚高的来源如果把这些被 abort 的请求计入 metric其 reward 全是 0相比不做 over sample 的 baselineover sample 的 reward 会显著偏低但如果完全无视它们又会引入另一种偏差——被 abort 的请求往往是轮数更多、难度更大的样本其 reward 天然低于未 abort 的请求无视这部分 reward 会导致 reward 虚高loss 侧的疑问直接让 loss mask 0 可以避免 padding 请求影响 loss但根据咨询结果这一做法可能只对特定的agg loss mode有效需要进一步研究agg_loss在不同模式下的行为。6.4 实测结论文档作者观点training step 的 reward 虚高与 6.3 的分析一致validation step 的 reward 能与 baseline 对齐作者认为 validation 阶段的 reward 是准确的目前并无大碍rollout time 有明显改善这正是引入 Over-Sample 的核心收益。七、Profiling 佐证被 abort 请求的耗时画像profile.md 记录了 Over-Sample 上线后对训练进行 profile 的方法与观测。为精确记录每个被取消请求的耗时需要把async_rollout_with_monitoring_start_time声明为nonlocal变量使其成为全局启动时间async def process_request_with_monitoring(req): nonlocal completed_count nonlocal async_rollout_with_monitoring_start_time try: result await self._async_rollout_a_request(req, do_sample, is_validate, **kwargs) async with completion_lock: if completed_count target_completion: completed_count 1 return result except asyncio.CancelledError: # request is cancelled, return padding logger.info(fRequest {req.request_id} was cancelled, creating padding) aborted_requests.append(req.request_id) torch.cuda.synchronize() req_aborted_time time.time() self.log_manager.log( self.log_path, eventaborted_request_with_cancelled_error, durationreq_aborted_time - async_rollout_with_monitoring_start_time, workidself._rank, stepself.step, extra{request_id: req.request_id}, ) self._create_padding_request(req) return except Exception as e: logger.error(fUncaught exception in process_request_with_monitoring: {e}) logger.error(This shall not happen, please check the code) raise e这样req_aborted_time - async_rollout_with_monitoring_start_time即为被 abort 请求含 abort 本身的总耗时。实际产出的 profiling 日志如下{timestamp: 2025-08-11T23:19:45.001830, event: aborted_request_with_cancelled_error_padding, duration_sec: 0.002585887908935547, extra: {request_id: 661a96a0-35e6-4662-9fff-bf9194bd3d49}, workid: 2, step: 4} {timestamp: 2025-08-11T23:19:45.001989, event: aborted_request_with_cancelled_error, duration_sec: 84.5104877948761, extra: {request_id: 45bb6a77-b6b2-4ed5-afde-97ccf622cdb0}, workid: 2, step: 4} {timestamp: 2025-08-11T23:19:45.004511, event: aborted_request_with_cancelled_error_padding, duration_sec: 0.0025205612182617188, extra: {request_id: 45bb6a77-b6b2-4ed5-afde-97ccf622cdb0}, workid: 2, step: 4} {timestamp: 2025-08-11T23:19:45.005685, event: async_rollout_with_monitoring_duration, duration_sec: 84.51417350769043, extra: {total_requests: 1024, target_completion: 921, completed_count: 921}, workid: 2, step: 4}这段日志透露了非常关键的信息该 worker 在该 step 发起total_requests 1024目标完成数target_completion 921约 90% 的完成率实际completed_count 921与设计完全吻合被 abort 的请求耗时约 84.5 秒与async_rollout_with_monitoring_duration84.51 秒几乎一致说明这些请求一直运行到截止时刻才被统一取消取消后的 padding 创建耗时仅约 2.5 毫秒可忽略不计两次aborted_request_with_cancelled_error_padding事件说明同一请求的 padding 创建与 abort 记录成对出现逻辑闭环。八、总结Over-Sample 的定位与后续方向将本指南的内容串起来可以得出以下结论动机明确多轮 rollout 的长尾效应80% 请求在前 40%50% 时间内完成严重拖慢训练Over-Sample 通过「超额发起 达到目标即截止」把等待最慢请求的时间转换为可控的固定开销实现精简monitor_and_cancel0.1 秒轮询 取消 engine abort、process_request_with_monitoring真实结果 or padding、run_with_cancellationgather(..., return_exceptionsTrue)统一收口三个函数即可覆盖完整流程架构上省事由于 verl 的AsyncEngine重写了abort_requestSGLang Engine 侧无需新增实现、无需发版副作用需持续治理padding 请求对 loss 的影响依赖loss mask 0的正确性且可能仅对特定 agg loss mode 有效reward 指标存在「计入则偏低、无视则虚高」的两难当前以 validation reward 作为可信基准GRPO group size 不一致对方差的影响尚无定量结论收益已被观测到rollout time 明显改善validation reward 能与 baseline 对齐。如果你希望进一步深入 verl 多轮 RL 的整体数据流与训练循环可继续阅读 readme-2.mdRayPPOTrainer.fit()与 make experience 全流程与 readme.md数据预处理、Hydra 分层配置若关心性能分析与长尾量化可参考 profile.md多轮 rollout 的整体背景与架构可参阅 verl-multiturn-rollout-Release.md。Over-Sample 特性本身仍处于「需要大改」的演进阶段其 padding 语义、reward 统计口径与 GRPO 方差影响是值得继续研究的方向。赞分享文档教程人工智能大模型RLHF【免费下载链接】Awesome-ML-SYS-TutorialMy learning notes for ML SYS.项目地址https://gitcode.com/gh_mirrors/aw/Awesome-ML-SYS-Tutorial点击查看免费下载相关推荐Ingress NGINX Controller 接入 OpenTelemetry分布式链路追踪配置与实战指南Ingress NGINX Controller 接入 OpenTelemetry分布式链路追踪配置与实战指南 导读 本文聚焦 Kubernetes 生态中最文档教程人工智能大模型RLHFcelld遥测完全指南OpenTelemetry trace写Parquet并用DuckDB查询celld遥测完全指南OpenTelemetry trace写Parquet并用DuckDB查询 celld 是一个自托管、分布式的 Durable Obje文档教程人工智能大模型RLHFverl 多轮Multi-turnRollout 完全指南SGLang 引擎、工具调用与增量 Tokenization 实战verl 多轮Multi turnRollout 完全指南SGLang 引擎、工具调用与增量 Tokenization 实战 导读 本文是 verlHy人工智能大模型强化学习RLHF分布式训练微调上一篇MoviePilot媒体库自动化连接异常3步诊断法与完整修复指南下一篇如何5分钟安装SD-PPPPhotoshop AI插件终极指南让你的设计效率提升300%创作声明:本文部分内容由AI辅助生成(AIGC),仅供参考
阅读完成 · 觉得有帮助?