Conversation
This file contains hidden or bidirectional Unicode text that may be interpreted or compiled differently than what appears below. To review, open the file in an editor that reveals hidden Unicode characters.
Learn more about bidirectional Unicode characters
Sign up for free
to join this conversation on GitHub.
Already have an account?
Sign in to comment
Add this suggestion to a batch that can be applied as a single commit.This suggestion is invalid because no changes were made to the code.Suggestions cannot be applied while the pull request is closed.Suggestions cannot be applied while viewing a subset of changes.Only one suggestion per line can be applied in a batch.Add this suggestion to a batch that can be applied as a single commit.Applying suggestions on deleted lines is not supported.You must change the existing code in this line in order to create a valid suggestion.Outdated suggestions cannot be applied.This suggestion has been applied or marked resolved.Suggestions cannot be applied from pending reviews.Suggestions cannot be applied on multi-line comments.Suggestions cannot be applied while the pull request is queued to merge.Suggestion cannot be applied right now. Please check back later.
PR type
PR information
reward_loop 奖励流水线:流式采样与节点微批(原理 / 收益 / 验证)
本 PR 的
reward_loop(src/twinkle/reward_loop/)把"算奖励"从 RL 训练主循环解耦成一条可异步提交 / 收集的流水线,并配套两项能力:引擎级流式采样(序列边生成、奖励边算)与奖励节点
微批(
BENCH_MICRO_BATCH,把逐条提交自动合并成大请求)。原理见"功能与原理"一节,实验证据见§1~§2。
收益总览
实测场景:RM 生成式 judge 打分,每步 32 条序列(数字为每步
t_collect,即采样结束后等待全部奖励就绪的耗时):
(§2.7)——整批必须等全部序列生成完,流式把奖励藏进了采样窗口。
run-路径-步骤语义检查点全 True(§1.2)。
功能与原理
reward_loop是 RL 训练中负责"算奖励"的独立组件:上层只做两件事——submit一批待打分的样本、collect取回分数。异步化的目的就是让奖励计算与采样、训练重叠(RewardLoopMetrics.overlap_ratio量化这个重叠),本报告回答的正是"这个重叠能换来多少墙钟收益、代价在哪里"。
数据流:采样出的序列 → 组装
RewardItem(item_id/solution_str/ground_truth/extra_info)→
pipeline.submit(items)返回BatchHandle→ worker 内的 manager 并发打分 →pipeline.collect(handle)按
item_id还原为与输入同序的RewardResult→ 交给 advantage 与训练。本报告对比的两条路径(整批 A / 流式 B)差别只在什么时候 submit,打分组件本身完全相同。
data.pyRewardItem/RewardResult;split_items按 worker 数分块,reorder_by_id+assemble_scores还原输入顺序(缺失/重复 item_id 直接报错)pipeline.pyAsyncRewardPipeline:submit/collect,backlog背压(满时阻塞或丢弃最旧)、on_backlog_full/on_error(报错或记 0 分);worker 是 Ray actor 时走ray.get,否则走内置线程池(宽度 =num_workers)worker.pyRewardLoopWorker:持有 manager,compute_score_batch(items)=asyncio.run(manager.run_batch(items));支持custom_reward_function_path动态导入打分函数reward_manager/RewardManagerBase(异步call_score/run_single/run_batch,normalize_score归一为标量分数,max_concurrent信号量)+ 注册表register/get_reward_manager_cls;内置naive/rate_limited(rpm / tpm / 超时兜底)/remote/dapo(超长惩罚)/gdpo,支持注册自定义 manager(本测的batch_judge即自定义)config.pyRewardLoopArgs:num_workers/manager_name/manager_source/backlog/mode/max_rpm/max_tpm/max_concurrent/timeout等metrics.pyRewardLoopMetrics:提交/收集批次数、submit / collect 耗时、max_backlog、overlap_ratio= 1 − collect 等待 /(submit + reward)default_score.pydata_source注册 scorer,未知来源按unknown_rewards告警或报错最小用法:
参考:
cookbook/rl/reward_loop/bench_streaming_local.py。能力一:引擎级流式采样(路径 B)
传统流程(路径 A)要等一步的全部序列生成完、算完奖励才能训练;流式(路径 B)用
vLLMSampler.sample_sequences_to_queue在一次远程调用内并发调度全部序列(vLLM 保持合批,不牺牲采样吞吐),任何一条序列先生成完就立刻回传事件、立刻提交打分。它的收益来源是重叠窗口:
序列完成得越分散(例如长短答案混合的长尾生成),越有机会在采样进行中把奖励算掉。正确性与整批
位级一致(§1.2),代价只是奖励侧需要合适的提交粒度(见能力二与 §2.4)。
能力二:奖励节点微批(
BENCH_MICRO_BATCH)流式的默认提交方式是逐条(保留逐条事件语义),但生成式 judge 的成本按请求次数计——一个请求
装 1 条与装 16 条耗时几乎相同(§2.4/§2.6),逐条提交等于把请求数放大 16 倍。微批在 worker 内把
多次
submit进来的 item 先攒着,攒满micro_batch_size条(或等micro_batch_timeout_ms兜底)后一次发给打分后端:上游提交方式完全不变,引擎侧请求数自动收敛。实现上 worker 通过
max_inflight_calls告知 pipeline 放大在飞请求数,否则逐条提交最多 2 个请求在飞、攒不满批。manager 侧只需覆写新增的
call_score_batch批接口(默认实现保持逐条语义,完全向后兼容)。收益:判词 200 token 时逐条
t_collect85.5s → 13.6s(§2.5),长尾场景反超整批(§2.7)。0. 口径与指标
vLLMSampler.sample_sequences_to_queue,一次远程调用内并发调度全部序列并逐条回传事件);提交粒度默认逐条,也可配为攒批后提交(见 §2)。
pipeline.submit装多少条——整批 = 全部序列;逐条 = 每完成一条就提交一条;每 2 条 = 攒够 2 条再提交。路径 A 恒为整批;路径 B 默认逐条。
奖励节点最终发给打分后端的"请求"大小还受 worker 内微批影响(§2.5)。
请求在跑。
submit的内容按 worker 数切出的分块(块数 = min(worker 数, 提交条数)),每块交给一个 worker。
结论之一是成本按请求计(§2.4)。
完成时刻);窗口越大,流式越有机会把奖励藏进采样时间里。
1. 函数式奖励(GSM8K 规则打分)
1.1 环境与配置
twinkle.initialize(mode='ray'),无 twinkle-server),昇腾 NPU ×4(ASCEND_RT_VISIBLE_DEVICES=2,3,4,5),model 1 卡 + sampler 1 卡ms://Qwen/Qwen3.5-4B(可覆盖为本地路径)messages+gold_answer),经DatasetMeta(data=rows)内存加载,完全离线gsm8k_score),毫秒量级(延迟=0 时t_collect≈ 0,见 §1.3);昂贵奖励由TWINKLE_REWARD_DELAY_MS注入TWINKLE_REWARD_NUM_WORKERS=2;路径 B 提交粒度per-item(固定)1.2 正确性
Level-1(确定性:temperature=0 + 固定 seed,不训练)
Level-2(语义:随机采样)
60 个 run-路径-步骤检查点全部 True。
1.3 收益测量
基础配置分解(batch=4, gen=4, max_tokens=1024, delay=0, 6 步均值)
单变量扫描(每 run 4 步;单位秒)
奖励尾部与重叠窗口
B 的首条奖励相对 A 的提前量(= A 的
reward_head_start− B 的reward_head_start,A 的该值≈其采样结束点):
机制(结论 3、4 的依据)
t_collect≈ delay(延迟=0 时 ≈0;d200 0.205~0.21s、d1000 1.005s,见上表):整批提交时单个 handle 内部由
run_batch用 asyncio 并发全部 item(2 个 chunk × 8 并发),奖励耗时几乎完全被并行吸收。
AsyncRewardPipeline线程池宽度(
max_workers = num_workers,本测为 2):d1000 下 16 条 ≈ (16/2) × 1s ≈ 8s,与实测最大尾部7.977s 吻合;由于同批序列完成时刻集中(引擎级流式合批),奖励请求成簇到达,尾部最明显。
max_concurrency=1的串行(
remote_class不传该参数即串行),此时提高 worker 数无效——加大提交 chunk 是通用解法。1.4 小结(函数式奖励)
rewards / advantages 位级等价,reward 一致性建立在非空 ground truth 上(非平凡)。
逐序列调用的串行代价——
t_sample与 A 同量级(20.3s vs 19.6s),7 组t_total的 B/A 落在1.00~1.19×(唯一超过 1.06× 的档位是 d1000),同时保留了逐条完成事件语义
(Level-1 / Level-2 检查通过)。
t_collect≈ 0),无可重叠成本;即使长输出下重叠窗口很大(tok2048 首条奖励提前 22.0s),
t_total仍持平。reward_head_start确实提前(gen8 5.5s、d200/d1000 约 4~5s),但
reward_tail_after_sample同步变大(B 最大 8.0s vs A 1.0s),净墙钟变差(d1000 B/A = 1.19×)。残余代价来自提交粒度而非采样方式;修法为整批/小批提交(或 §2.5 的奖励节点攒批)。
2. 奖励模型(RM,生成式 judge)
2.1 配置
DeviceGroup('reward')单卡(TWINKLE_REWARD_MODEL_ID,默认跟随TWINKLE_MODEL_ID)vLLMSampler(remote_group='reward');compute_score四参数契约;输出Correct/Incorrect或 1/0 → 1.0/0.0AsyncRewardPipeline按 chunk 并行调用 judge(多线程 → vLLM 请求合批),与"并行 API 型 RM"一致TWINKLE_REWARD_JUDGE_MAX_TOKENS=8(判词 8 token)判词 = judge 为每条样本生成的判定文本(含推理与结论),长度由
TWINKLE_REWARD_JUDGE_MAX_TOKENS限制;判词长度直接决定单请求成本(§2.6)。
2.2 提交粒度矩阵(取稳定步均值;首步含 judge 引擎一次性 warm-up,已剔除)
结论(RM 场景)
t_collect0.57~0.94s),引擎对并发请求合批(探针 C≈A,§2.6),三档
t_total11.7~12.2s 等价。对照 §1.3 的 d1000(函数式奖励 + 注入 1s/条延迟、每步 16 条)中同样的逐条提交:尾部被放大到 7.9s(A 侧 1.0s),
说明粒度差异只在"奖励贵到超出重叠窗口"时才出现。
t_collect0.78s vs 0.57s;A 的t_total偏高来自首步 warm-up 与采样侧(§4)。
reward_mean=0.0000且无解析失败告警——judge 全部判 Incorrect(4B 验证器判别失效)。延迟特征(真实推理耗时)有效,判别语义无效;性能对比结论不受影响,
如需真实判别需换更强 judge 或判别式 RM(直接打分、不生成判词,因此也没有判词长度问题)。
本地 harness 下可用并发 = pipeline 线程池宽度 =
num_workers= 2)相对可重叠窗口的大小,而非单看单条延迟:
8.5~9.3s → 三档无差别,粒度无关紧要;
(首条奖励早于采样结束约 4.2s,见 §1.3)→ 实测尾部溢出到 1.2~8.0s,此时应改用整批/小批提交;
收益越明显(与排队溢出无关,故其收益在串行总量不溢出时也会出现)。
2.3 边界验证:长尾生成 + 昂贵 judge(采样 max_tokens=2048 / judge 判词 200 token)
配置同 §2.1 的矩阵(每步 4 条序列)。各臂 3 步;B 各臂的最慢一步是 step2(该批序列生成偏长,
t_sample32.5~36.1s),A 的最慢一步是 step0(首步 warm-up),故t_total同时给出剔除最慢一步后的值。
结论(边界)
t_collect5.9~7.9s,仍低于采样 24~28s;流式未因奖励变贵而反超。B 的
t_total整体小于 A(31.5~36.9s vs 40.6s),差异来自采样侧(A 的非流式采样 27.9s vs B约 24~27s)与 A 的首步 warm-up,而非奖励侧。
§2.4:整批请求数最少;小批(每 2 条)在 2 个 worker 下被切成 1 条/块,实际 chunk 与逐条相同。
semantic_ok列不可用(reward_deterministic检查拿 judge 分数与规则奖励比对,语义不成立),应忽略。奖励时间线已给
batch_judge补上逐条reward_start/reward_end与每 chunk 的
judge_chunk/judge_call事件,因此 §2.4 的排队与合批指标同样适用于 RM。2.4 RM 压力档:32 条/步(判词 200 token)
配置:batch=4 × gen=8(32 条/步)、采样 max_tokens=512、判词 200 token
(
TWINKLE_REWARD_JUDGE_MAX_TOKENS=200)、3 步;奖励并发TWINKLE_REWARD_NUM_WORKERS=2;其余同 §2.1。同一配置把并发改成 4(
TWINKLE_REWARD_NUM_WORKERS=4)另跑一次作对照(逐条臂):t_collect85.5s → 45.8s、chunk 级并发 4——提高并发有效(约线性)。
该跑采样侧与训练侧正常(
t_sample11.0~13.0s、t_train6.9~7.7s),差异只在奖励侧。结论(压力档)
1.6×)。A 与 B 整批的奖励侧一致(9.38s vs 8.93s)。
请求只贵 1.7×(vLLM 合批,探针 C≈A,§2.6);32 个请求与 2 个请求的数量差(16×)才是主因,
总耗时由请求个数决定。
t_collect85.5s → 45.8s(≈线性,因为引擎合批并发请求);微批放大在飞请求数同样有效(§2.5)。
AsyncRewardPipeline按 worker 数轮询切分提交内容,实际 judgechunk = 提交条数 ÷ worker 数(2 worker → 16 条/块,4 worker → 8 条/块)。因此"每 2 条一批"在
2 或 4 个 worker 下都会被切成 1 条/块,退化成逐条(表中 mini 与 per-item 的调用次数与平均条数
完全相同)。要形成真正的小批,提交条数必须大于 worker 数(
BENCH_MINI_BATCH_SIZE>TWINKLE_REWARD_NUM_WORKERS)。整体落在采样之后,B 相对 A 没有提前量可用——与 §1.3 现象一致(该结论仅对短输出成立;长尾下
窗口真实存在,见 §2.7)。
2.5 奖励节点攒批(
BENCH_MICRO_BATCH):上游逐条提交 + 引擎大请求微批定义见"功能与原理 / 能力二"。配置同 §2.4(32 条/步、判词 200 token、
TWINKLE_REWARD_NUM_WORKERS=2),微批为BENCH_MICRO_BATCH=16/BENCH_MICRO_BATCH_TIMEOUT_MS=200(攒满 16 条或等 200ms 即发起一次 judge 调用);"关"一行为逐条基线。另做
BENCH_MICRO_BATCH大小扫描,判词长度扫描见 §2.6。
微批与引擎侧修复的关系(重要):微批是为绕过 Ray actor 阻塞设计的调用方侧方案——早期训练机
状态上,同步阻塞的
sample占住 asyncio actor 的事件循环,并发请求在进引擎前被串行化(探针4 并发 = 4 × 单请求),逐条提交因此付出 9.6× 尾部代价。main 分支 commit
c839a4e(modelscope#280)已在引擎侧修复该问题:
vLLMSampler.sample声明enable_continous_work=True,async companion 把阻塞体挪进线程池,并发请求得以一起进引擎合批——本次探针已证实其在训练机生效(§2.6:C ≈ A)。
因此:在该修复生效的版本上,逐条提交的尾部随在飞请求数下降(workers 2→4:85.5s → 45.8s);
调大 worker 数或后续引擎优化都可能进一步替代微批。微批的价值在于:不增加 worker 数、不改上游
提交方式,即把尾部压到接近整批(13.6s vs 8.9s),并保留逐条事件语义。
微批大小扫描(逐条臂、
BENCH_MICRO_BATCH_TIMEOUT_MS=200)扫描只跑逐条臂(
BENCH_RUNS=rm-b-per);末行"整批(对照)"为同配置的 rm-b-whole,作参考线。结论(攒批)
t_collect85.5s → 13.6s(6.3×);仍差整批 4.7s——整批的 2 个 16 条请求被引擎合批(每调用 8.85s),逐条的块按到达节奏
串行(块 {1,15,16},未攒满)。
先于攒满触发。大小扫描单调改善(4→47.7s、8→23.5s、16→13.6s),MB=16 已接近整批,无需更大。
提高 worker 数等效(§2.4 结论 3),微批(
BENCH_MICRO_BATCH)与TWINKLE_REWARD_NUM_WORKERS都能放大在飞请求数。
mini与微批组合不再差(14.2s ≈ 逐条 13.6s):引擎合批后小块的代价消失,旧结论反转;三档提交方式 + 微批都收敛到整批量级。
2.6 判词长度对 judge 调用成本的影响
探针(同 prompt 三形态:1 调用 1 条 / 1 调用 4 条 / 4 并发各 1 条)与压力档 bench(MB=16)
联合测量:
"端到端"列的"整批"即 §2.5 的 rm-b-whole(非流式提交),作对照线;探针三列(1×1 / 1×4 /
4 并发×1)是引擎层单请求测量,与提交路径(流式/整批)无关。
结论(判词长度)
预填充固定 ≈0.3s。
500 token 处自然 EOS,判词上限再大也不再增加成本("不限制"最多就是这个量级,不会失控到分钟级)。
并发请求被合批处理(早期版本的"跨请求串行 C=4×A"与旧代码有关,§4)。
26.7)——判词越长整批优势越大(2 个并发 16 条请求被合批,逐条块按到达节奏串行);判词上限决定
绝对成本(8→800 约 15×)。
2.7 长尾 × 微批:流式首次反超整批
配置同 §2.4 但采样
max_tokens=2048(长尾生成);判词 200 token、32 条/步、2 worker。MB=0 为基线,MB=16 分 200ms / 1000ms 两档超时。
MB=0 跑的整批臂 judge 调用偏慢(31.3s,首跑 warm-up),整批对照取后两跑(11.4~11.5s)。
结论(长尾 × 微批)
head17~20s 远小于t_sample41~43s,奖励在采样进行中就开始计算(§2.4 结论 5 仅对短输出成立)。
t_collect7.10s < 整批 11.49s,t_total58.6s < 69.0s(快 ≈10s)——整批必须等全部序列完成才开始(
head= 41s),流式 + 微批把奖励藏进了采样窗口。200ms / 1000ms 都攒不满 16),真正起作用的是微批把在飞请求数从 2 放大到 10~16(
conc),引擎对并发请求合批(探针 C≈A,§2.6);配合长尾的真实重叠窗口(结论 1),奖励被藏进采样时间。
3. 总体结论
RM 模式需忽略
semantic_ok(检查语义不适用于 judge)。超过 1.06× 的档位是 d1000);RM 短判词下持平;长判词下整批提交更优(§2.4);长尾 × 微批下流式
逐条反超整批(§2.7)。
尾部等待抵消提前开始的收益。扫描表也没有出现收益随变量增大而放大的趋势——delay 0→200→1000ms
的 B/A 为 1.01→1.06→1.19×,max_tokens 512→1024→2048 为 1.06→1.01→1.01×,gen 2→4→8 为
1.02→1.01→1.00×。长尾档(§2.3,2048 token)B 的
t_total整体低于 A(31.5~36.9s vs 40.6s),但该差异在采样侧(A 的非流式采样更长、首步 warm-up),且 §1 同长度档(tok2048)A ≈ B
(36.9s vs 36.8s),故不归为流式收益。真正的例外是长尾 × 微批(§2.7):序列完成分散带来
真实重叠窗口,且引擎对并发请求合批,逐条 + 微批的
t_collect7.10s 反超整批 11.49s、t_total快 ≈10s——这是唯一一次流式获得实质墙钟收益。对并发请求合批(探针 C≈A,§2.6),但逐条提交受 pipeline 池宽限制(无微批时在飞请求数 =
num_workers)——提高 worker 数(§2.4:2→4 使逐条 85.5s → 45.8s)或开启微批放大在飞数(§2.5:85.5s → 13.6s)都有效;短输出下整批仍最优(8.9s)。注意小批要真正生效,提交条数必须
大于 worker 数(§2.4 结论 4)。
4. 数据可信度与边界
噪声量级(组内步间抖动可达 ±2%)。
reward_tail_after_sample的小负值(如基础配置 B 侧的 −0.001s)表示最后一条奖励在采样结束打点前就绪,不是计时错误。
单请求 ≈0.3s,200 token 单请求 ≈5s,且上限超过约 500 token 后由 judge 的自然 EOS 封顶
(§2.6)。跨批次比较前先核对
TWINKLE_REWARD_JUDGE_MAX_TOKENS(bench 启动日志会打印judge_max_tokens=)。t_sample/t_train,不影响 §2 的奖励侧结论。(C ≈ A;修复为 main 分支 commit
c839a4e的 async companion /enable_continous_work)。本文 §2 数字为当前代码重测所得,跨版本比较前先核对
TWINKLE_SAMPLER_MAX_CONCURRENCY与 async-companion 状态(探针头部会打印)。