Skip to content

[feature]reward_loop - #284

Open
doctorMcy wants to merge 1 commit into
modelscope:mainfrom
doctorMcy:feature_rl_reward_loop
Open

doctorMcy wants to merge 1 commit into
modelscope:mainfrom
doctorMcy:feature_rl_reward_loop

Conversation

@doctorMcy

Copy link
Copy Markdown

PR type

  • Bug Fix
  • [ √ ] New Feature
  • Document Updates
  • More Models or Datasets Support

PR information

奖励测试数据与结论(函数式奖励 / 奖励模型)

功能介绍:函数式奖励(GSM8K 规则打分,毫秒量级)下,引擎级流式采样与整批采样端到端持平
(两次运行 B/A = 0.99~1.06×;注入 1s 奖励延迟的 d1000 档为 1.19~1.20×);真实奖励模型(生成式
judge)下两条采样路径同样持平,提交粒度在奖励侧串行量远小于采样窗口时无关紧要,长判词下整批最优。
两级正确性检查全部通过。

背景:reward loop 概览

reward_loopsrc/twinkle/reward_loop/)是 RL 训练中负责"算奖励"的独立组件:把打分从训练主循环里
解耦成一条可异步提交 / 收集的流水线,上层只做两件事——submit 一批待打分的样本、collect 取回
分数。异步化的目的就是让奖励计算与采样、训练重叠(RewardLoopMetrics.overlap_ratio 量化这个重叠),
本报告回答的正是"这个重叠能换来多少墙钟收益、代价在哪里"。

数据流:采样出的序列 → 组装 RewardItemitem_id / solution_str / ground_truth / extra_info
pipeline.submit(items) 返回 BatchHandle → worker 内的 manager 并发打分 → pipeline.collect(handle)
item_id 还原为与输入同序RewardResult → 交给 advantage 与训练。本报告对比的两条路径
(整批 A / 流式 B)差别只在什么时候 submit,打分组件本身完全相同。

模块 职责
data.py 契约与工具:RewardItem / RewardResultsplit_items 按 worker 数分块,reorder_by_id + assemble_scores 还原输入顺序(缺失/重复 item_id 直接报错)
pipeline.py AsyncRewardPipelinesubmit / collectbacklog 背压(满时阻塞或丢弃最旧)、on_backlog_full / on_error(报错或记 0 分);worker 是 Ray actor 时走 ray.get,否则走内置线程池(宽度 = num_workers
worker.py RewardLoopWorker:持有 manager,compute_score_batch(items) = asyncio.run(manager.run_batch(items));支持 custom_reward_function_path 动态导入打分函数
reward_manager/ RewardManagerBase(异步 call_score / run_single / run_batchnormalize_score 归一为标量分数,max_concurrent 信号量)+ 注册表 register / get_reward_manager_cls;内置 naive / rate_limited(rpm / tpm / 超时兜底)/ remote / dapo(超长惩罚)/ gdpo,支持注册自定义 manager(本测的 batch_judge 即自定义)
config.py RewardLoopArgsnum_workers / manager_name / manager_source / backlog / mode / max_rpm / max_tpm / max_concurrent / timeout
metrics.py RewardLoopMetrics:提交/收集批次数、submit / collect 耗时、max_backlogoverlap_ratio = 1 − collect 等待 /(submit + reward)
default_score.py 默认规则打分分派:按 data_source 注册 scorer,未知来源按 unknown_rewards 告警或报错

0. 口径与指标

  • 路径 A(整批):整批采样 → 整批提交奖励 → collect → 训练。
  • 路径 B(流式):引擎级逐序列流式采样(vLLMSampler.sample_sequences_to_queue,一次远程调用内
    并发调度全部序列并逐条回传事件),任一序列完成即提交该条奖励。
指标 定义(单位:秒)
t_sample 采样阶段墙钟
t_submit 提交阶段墙钟(本测中可忽略,表内省略)
t_collect 等待本步全部奖励就绪的墙钟(奖励尾部等待)
t_train forward_backward + clip_grad_and_step 墙钟
t_total 单步总墙钟 = 四者之和
seq_per_s 每步序列数 / t_total
reward_head_start 本步第一条奖励开始计算距采样开始的时长(B 中 < t_sample 表示与采样重叠)
reward_tail_after_sample 本步最后一条奖励就绪距采样结束的时长(负值 = 早于采样结束就绪)

1. 函数式奖励(GSM8K 规则打分)

1.1 环境与配置

机器 / 部署 训练机,本地 ray 模式(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(可覆盖为本地路径)
数据集 GSM8K train(7473 行),本地 jsonl(messages + gold_answer),经 DatasetMeta(data=rows) 内存加载,完全离线
采样参数 temperature=1.0, top_p=0.95, logprobs=1, num_samples=1
每步序列数 batch4 × gen(base/tok/d200/d1000=16,gen2=8,gen8=32)
奖励 函数式规则打分(gsm8k_score),毫秒量级(延迟=0 时 t_collect ≈ 0,见 §1.3);昂贵奖励由 TWINKLE_REWARD_DELAY_MS 注入
奖励并发 TWINKLE_REWARD_NUM_WORKERS=2;路径 B 提交粒度 per-item(默认)
调度 单缓冲:本步采样完成 → collect 本步奖励 → 训练(避免双缓冲掩盖流式收益)

1.2 正确性

Level-1(确定性:temperature=0 + 固定 seed,不训练)

检查项 最终结果
样本总数一致(16=16)
tokens 逐位一致
logprobs 逐位一致 ✅(max diff 0.0)
decoded 文本一致
reward 逐样本一致
优势(GRPOAdvantage)一致

Level-2(语义:随机采样)

检查项 A(批量) B(流式)
无丢样本 / 无重复 item_id
handle 结果与 item_id 对齐
reward 函数确定性(重打分一致)
优势组均值归零

60 个 run-路径-步骤检查点全部 True(两次独立运行均如此)。

1.3 收益测量

基础配置分解(运行 1:batch=4, gen=4, max_tokens=1024, delay=0, 6 步均值)

指标 A(批量) B(引擎流式)
t_sample 19.6 20.3
t_collect 0.00 0.00
t_train 3.7 3.6
t_total 24.4 24.5
seq/s 0.65 0.65
reward_head_start 19.7 ≈16
reward_tail_after_sample 0.00 -0.001

单变量扫描(每 run 4 步;单位秒;两次独立运行同配置、同 harness,用于互相复核)

run 变量 运行 1 A 运行 1 B 运行 1 B/A 运行 2 A 运行 2 B 运行 2 B/A
base 24.4 24.5 1.00× 24.35 24.64 1.01×
gen2 gen=2 21.2 21.2 1.00× 20.90 21.31 1.02×
gen8 gen=8 32.7 33.0 1.01× 32.74 32.73 1.00×
tok512 max_tokens=512 17.4 18.2 1.05× 16.00 17.02 1.06×
tok2048 max_tokens=2048 44.6 44.1 0.99× 44.98 45.50 1.01×
d200 delay=200ms 25.5 26.9 1.05× 25.16 26.56 1.06×
d1000 delay=1000ms 26.3 31.5 1.20× 25.98 30.96 1.19×

奖励尾部与重叠窗口

run A t_collect(运行 1 / 运行 2) B t_collect(运行 1 / 运行 2)
d200 0.21 / 0.205 ≈0.4~1.6 / 0.397~1.578
d1000 1.005 / 1.005 ≈1.7~7.9 / 1.208~7.977

运行 2 下 B 的首条奖励相对 A 的提前量(= A 的 reward_head_start − B 的 reward_head_start
A 的该值≈其采样结束点):

run base gen2 gen8 tok512 tok2048 d200 d1000
提前量(s) 1.2 1.3 5.5 -0.9(更晚) 22.0 4.6 4.2

机制(结论 3、4 的依据)

  • A 侧 t_collect ≈ delay(延迟=0 时 ≈0;d200 0.205~0.21s、d1000 1.005s,见上表):整批提交时
    单个 handle 内部由 run_batch 用 asyncio 并发全部 item(2 个 chunk × 8 并发),奖励耗时几乎
    完全被并行吸收。
  • B 侧逐条提交时每次调用的 item 并发为 1,实际并发上限 = AsyncRewardPipeline 线程池宽度
    max_workers = num_workers,本测为 2):d1000 下 16 条 ≈ (16/2) × 1s ≈ 8s,与实测最大尾部
    7.977s 吻合;由于同批序列完成时刻集中(引擎级流式合批),奖励请求成簇到达,尾部最明显。
  • 在 Ray actor 部署(server 侧)中,逐条提交会撞上 actor 默认 max_concurrency=1 的串行
    remote_class 不传该参数即串行),此时提高 worker 数无效——加大提交 chunk 是通用解法

1.4 小结(函数式奖励)

  1. 正确性:两级检查全部通过;Level-1 在 greedy + 固定 seed 下达成 tokens / logprobs / decoded /
    rewards / advantages 位级等价,reward 一致性建立在非空 ground truth 上。
  2. 采样层面持平:路径 B 在一次远程调用内并发调度全部序列(vLLM 保持合批),没有引入
    逐序列调用的串行代价——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 检查通过)。
  3. 延迟=0 时流式重叠买不到墙钟:奖励本身耗时可忽略(延迟=0 时 A/B 的 t_collect ≈ 0),
    无可重叠成本;即使长输出下重叠窗口很大(tok2048 首条奖励提前 22.0s),t_total 仍持平。
  4. 延迟>0 时奖励侧吞吐成为瓶颈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 配置

打分器 生成式 judge,独立 DeviceGroup('reward') 单卡(TWINKLE_REWARD_MODEL_ID,默认跟随 TWINKLE_MODEL_ID
judge 形态 冻结权重的第二个 vLLMSamplerremote_group='reward');compute_score 四参数契约;输出 Correct/Incorrect 或 1/0 → 1.0/0.0
并发形态 AsyncRewardPipeline 按 chunk 并行调用 judge(多线程 → vLLM 请求合批),与"并行 API 型 RM"一致
矩阵 batch=2 × gen=2(每步 4 条)、max_tokens=512、3 步、TWINKLE_REWARD_JUDGE_MAX_TOKENS=8(判词 8 token)
粒度档 BENCH_SUBMIT_GRANULARITY=whole(整批 1 handle)/ mini(每 2 条一批)/ per-item(逐条)

2.2 提交粒度矩阵(取稳定步均值;首步含 judge 引擎一次性 warm-up,已剔除)

run 路径 粒度 t_sample t_collect t_total
rm-whole A 整批(1 handle) 8.5s 0.96s 11.3s
rm-b-whole B 整批 9.2s 0.97s 11.6s
rm-b-mini B 每 2 条一批 9.2s 1.30s 12.0s
rm-b-per B 逐条 9.3s 1.27s 11.9s

结论(RM 场景)

  1. 提交粒度差异从"数量级"缩至"噪声级":真实 judge 评分在亚秒级(8 token 判词;实测
    t_collect 0.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 开销),三档实际可视为等价。
  2. 采样路径(A vs B)在 RM 场景同样持平(11.3s vs 11.6s;t_sample 8.5s vs 9.2s)。
  3. ⚠️ judge 判别质量限制reward_mean=0.0000 且无解析失败告警——judge 全部判 Incorrect
    (4B 验证器判别失效)。延迟特征(真实推理耗时)有效,判别语义无效;性能对比结论不受影响,
    如需真实判别需换更强 judge 或判别式 RM 权重。
  4. 实践含义:是否需要调整提交粒度,取决于奖励侧的串行总量(条数 × 单条延迟 ÷ 可用并发,
    本地 harness 下可用并发 = pipeline 线程池宽度 = num_workers = 2)相对可重叠窗口的大小,
    而非单看单条延迟:
    • 本矩阵:每步 4 条、judge 单条在亚秒级(≲1s)→ 串行总量 ≲ 4 × 1s ÷ 2 = 2s,远小于采样窗口
      8.5~9.3s → 三档无差别,粒度无关紧要;
    • §1.3 的 d1000:每步 16 条、注入 1s/条 → 串行总量 ≈ 16 × 1s ÷ 2 = 8s,超过可重叠窗口
      (首条奖励早于采样结束约 4.2s,见 §1.3)→ 实测尾部溢出到 1.2~8.0s,此时应改用整批/小批提交;
    • §2.3 的长判词边界是第二条独立机制:chunk 越大,单次调用内并发的条数越多,judge 侧合批
      收益越明显(与排队溢出无关,故其收益在串行总量不溢出时也会出现)。

2.3 边界验证:长尾生成 + 昂贵 judge(采样 max_tokens=2048 / judge 判词 200 token)

run t_sample t_collect t_total(含 step0 抖动)
rm-whole / A 27.8s 12.1s 46.7s(剔除 step0 后 ≈41.5s)
rm-b-whole / B 25.1s 11.96s 39.0s(剔除 step0 后 ≈40.8s)
rm-b-mini / B 25.2s 17.4s 44.3s
rm-b-per / B 24.9s 15.1s 41.7s

结论(边界)

  1. 奖励成为关键路径的幅度不足:judge 合批后 t_collect 11.96~17.4s,仍低于采样 25~28s——
    流式(B)与批量(A)持平(t_total 差 ≈1s,噪声级),流式未因奖励变贵而反超
  2. 提交粒度首次显著分化:整批提交(t_collect 11.96s)优于小批(17.4s,-31%)与逐条
    (15.1s,-21%)——长判词下 judge 推理合批收益真实显现(4 条共享 prefill/decode)。注意这里的
    机制与 §1.3 的 d1000 不同:d1000 的瓶颈是逐条提交把并发压到线程池宽度后的排队等待,而本 run
    的串行总量并未溢出采样窗口,收益来自 judge 侧合批(chunk 越大,单次调用内可并发的 item 越多)。
  3. RM 模式下 semantic_ok 列不可用reward_deterministic 检查拿 judge 分数与规则奖励比对,
    语义不成立),应忽略;reward_head_start / reward_tail_after_sample 同样为空
    (batch_judge manager 未打奖励时间线点)。

3. 总体结论

  1. 正确性:函数式奖励与 RM 两条路径均通过(Level-1 位级等价;Level-2 60 个检查点全 True)。
    RM 模式需忽略 semantic_ok(检查语义不适用于 judge)。
  2. 性能持平:两种采样路径在函数式奖励(两次运行的 B/A 分别为 0.99~1.20× 与 1.00~1.19×,唯一
    超过 1.06× 的档位是 d1000)与 RM(11.3s vs 11.6s)下均持平;B 在持平的同时保留了逐条完成事件
    语义。
  3. 流式不带来墙钟收益(当前调度下):延迟=0 时无可重叠成本;延迟>0 时奖励侧吞吐成为新瓶颈,
    尾部等待抵消提前开始的收益。扫描表也没有出现收益随变量增大而放大的趋势——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 的比值)。
  4. 提交粒度是最有效的杠杆:奖励便宜(毫秒量级/条)时粒度无关;奖励昂贵且请求成簇到达时,
    逐条提交的并发上限 = pipeline 线程池宽度(本地)或 actor 串行度(Ray 部署),应改用整批/小批提交。
  5. 推荐用法:本地 ray 模式用 sample_sequences_to_queue(与整批性能持平 + 保留逐条事件语义);
    奖励端按成本选择提交粒度;server 侧(actor 部署)避免逐条提交。

Sign up for free to join this conversation on GitHub. Already have an account? Sign in to comment

Labels

None yet

Projects

None yet

Development

Successfully merging this pull request may close these issues.

1 participant