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
奖励测试数据与结论(函数式奖励 / 奖励模型)
功能介绍:函数式奖励(GSM8K 规则打分,毫秒量级)下,引擎级流式采样与整批采样端到端持平
(两次运行 B/A = 0.99~1.06×;注入 1s 奖励延迟的 d1000 档为 1.19~1.20×);真实奖励模型(生成式
judge)下两条采样路径同样持平,提交粒度在奖励侧串行量远小于采样窗口时无关紧要,长判词下整批最优。
两级正确性检查全部通过。
背景:reward loop 概览
reward_loop(src/twinkle/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告警或报错0. 口径与指标
vLLMSampler.sample_sequences_to_queue,一次远程调用内并发调度全部序列并逐条回传事件),任一序列完成即提交该条奖励。
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 收益测量
基础配置分解(运行 1:batch=4, gen=4, max_tokens=1024, delay=0, 6 步均值)
单变量扫描(每 run 4 步;单位秒;两次独立运行同配置、同 harness,用于互相复核)
奖励尾部与重叠窗口
运行 2 下 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 两次运行分别落在 0.99~1.20× 与 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~1.20×)。残余代价来自提交粒度而非采样方式;修法为整批/小批提交(
BENCH_SUBMIT_GRANULARITY=whole|mini)。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)BENCH_SUBMIT_GRANULARITY=whole(整批 1 handle)/mini(每 2 条一批)/per-item(逐条)2.2 提交粒度矩阵(取稳定步均值;首步含 judge 引擎一次性 warm-up,已剔除)
结论(RM 场景)
t_collect0.96~1.30s 为等到本批最后一条的耗时,单条不超过该值),串行代价基本被 2 路并发吸收——
t_collect差 ≈30%、t_total差 ≈6%。对照 §1.3 的 d1000(函数式奖励 + 注入 1s/条延迟、每步 16 条)中同样的逐条提交:
t_collect被放大到 1.2~8.0s(整批侧为 1.005s)。whole 略优(judge 合批 +无 handle 开销),三档实际可视为等价。
t_sample8.5s vs 9.2s)。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)
结论(边界)
t_collect11.96~17.4s,仍低于采样 25~28s——流式(B)与批量(A)持平(
t_total差 ≈1s,噪声级),流式未因奖励变贵而反超。t_collect11.96s)优于小批(17.4s,-31%)与逐条(15.1s,-21%)——长判词下 judge 推理合批收益真实显现(4 条共享 prefill/decode)。注意这里的
机制与 §1.3 的 d1000 不同:d1000 的瓶颈是逐条提交把并发压到线程池宽度后的排队等待,而本 run
的串行总量并未溢出采样窗口,收益来自 judge 侧合批(chunk 越大,单次调用内可并发的 item 越多)。
semantic_ok列不可用(reward_deterministic检查拿 judge 分数与规则奖励比对,语义不成立),应忽略;
reward_head_start/reward_tail_after_sample同样为空(batch_judge manager 未打奖励时间线点)。
3. 总体结论
RM 模式需忽略
semantic_ok(检查语义不适用于 judge)。超过 1.06× 的档位是 d1000)与 RM(11.3s vs 11.6s)下均持平;B 在持平的同时保留了逐条完成事件
语义。
尾部等待抵消提前开始的收益。扫描表也没有出现收益随变量增大而放大的趋势——delay 0→200→1000ms
的 B/A 为 1.00→1.05→1.20×,max_tokens 512→1024→2048 为 1.05→1.00→0.99×,gen 2→4→8 为
1.00→1.00→1.01×(均为运行 1 的比值)。
逐条提交的并发上限 = pipeline 线程池宽度(本地)或 actor 串行度(Ray 部署),应改用整批/小批提交。
sample_sequences_to_queue(与整批性能持平 + 保留逐条事件语义);奖励端按成本选择提交粒度;server 侧(actor 部署)避免逐条提交。