
AReaL 指标跟踪系统深度解析从流式 Rollout 到批量训练的统一统计管线【免费下载链接】AReaLThe RL Bridge for LLM-based Agent Applications. Made Simple Flexible.项目地址: https://gitcode.com/GitHub_Trending/are/AReaL本篇技术指南围绕 AReaL 仓库的指标跟踪Metrics Tracking机制展开系统讲解以areal.utils.stats_tracker为核心的统一指标收集、分布式归约与多后端日志落盘方案。文章将带你掌握「流式指标 批量指标」双范式的工作方式、scalar()/denominator()/stat()三种记录 API 的适用场景、层级作用域与命名跟踪器的用法以及如何通过stats_logger配置把指标送入 wandb、SwanLab、TensorBoard 等实验跟踪后端。读完即可在自定义 Rollout 工作流或训练引擎中正确埋点并理解每一步指标从收集、聚合到提交的完整链路。一、指标跟踪系统在 AReaL 中的定位AReaL 是面向 LLM Agent 应用的强化学习训练框架一次训练运行同时存在两类分布式的统计来源Rollout 工作器异步执行工作流如 RLVRWorkflow每个工作流独立完成一次采样与奖励计算完成时间天然不等训练引擎在数据并行DPrank 之间同步处理批次例如 PPOActor 的ppo_update()需要保证所有 rank 统计结果一致。这两类来源的同步模型完全不同因此 AReaL 提供了两种针对各自场景优化的指标范式范式适用对象记录方式聚合时机流式指标Rollout 工作器每个工作流完成时追加一个标量导出时由控制器做加权平均批量指标训练引擎带布尔掩码的全批次张量导出时跨数据并行组 all-reduce整个体系构建在 areal/utils/stats_tracker.py 的DistributedStatsTracker之上由 areal/utils/stats_logger.py 的StatsLogger负责最终提交到外部日志后端。二、核心组件DistributedStatsTrackerDistributedStatsTracker是线程安全的统计收集器通过Lock保护内部状态提供了四类核心能力命名跟踪器Named Trackers不同组件拥有隔离的指标命名空间层级作用域Hierarchical Scoping把相关指标组织成逻辑分组导出时自动拼出group/sub/key形式的键分布式聚合Distributed Aggregation导出时根据传入的reduce_group跨工作器自动归约多种归约类型Reduce Types平均值、求和、最小/最大值、标量加权平均。模块底部暴露了模块级便捷函数见 stats_tracker.py因此可以直接这样使用from areal.utils import stats_tracker # 默认跟踪器训练指标 stats_tracker.scalar(learning_rate0.001) # 命名跟踪器Rollout 指标 stats_tracker.get(rollout).scalar(reward0.5)get(name)会按名称创建/复用独立的DistributedStatsTracker实例默认名称对应训练指标跟踪器export_all()则遍历所有已注册的跟踪器并合并结果若跨 rank 出现重复键会发出告警见 stats_tracker.py。三、两种日志范式详解3.1 流式指标面向 Rollout 工作器Rollout 工作器异步执行工作流每个工作流在完成时单独记录标量指标累积在工作器进程内部列表中记录期间工作器之间无需任何同步归约全部推迟到导出阶段由控制器完成。以 RLVRWorkflow._collect_samples 为真实示例async def _collect_samples(self, engine, req, prompt_str, task_data): resp await engine.agenerate(req) reward await self._compute_rewards(resp, prompt_str, task_data) # 记录单个标量 - 追加到内部列表 # workflow_context.stat_scope() 自动区分评估/训练作用域 stats_tracker.get(workflow_context.stat_scope()).scalar(rewardreward) return resp, reward这里workflow_context.stat_scope()实现在 areal/infra/workflow_context.py根据当前上下文返回eval-rollout评估模式或rollout训练模式从而把评估与训练的 Rollout 指标自动隔离到不同命名空间。在自定义工作流中可以记录任意其他标量例如交互轮数、最大 token 数等async def run(self, data, **extra_kwargs): # workflow_context.stat_scope() 自动区分评估/训练作用域 stats_tracker.get(workflow_context.stat_scope()).scalar( num_turnsnum_turns, max_tokensmax_tokens, rewardreward ) return reward控制器聚合RolloutController.export_stats()见 areal/infra/controller/rollout_controller.py通过_collective_rpc从所有工作器收集各自已经完成本地归约的export_stats结果再由_merge_worker_stats合并标量均值按__count计数做加权平均base/avg|min|max形式的分布类指标保留加权均值与极值PRM 的 count/sum 类指标按 SUM 语义直接累加见 rollout_controller.py。这样即使各工作器处理的工作流数量不等、完成时间不一最终均值也能按真实样本数加权而不是简单求算术平均。3.2 批量指标面向训练引擎训练引擎在数据并行 rank 之间同步处理批次。指标以带布尔掩码分母的全批次张量记录导出时执行 all-reduce保证各 rank 拿到完全一致的统计结果。以 PPOActor.ppo_update 中的记录逻辑为例与官方文档给出的模式一致def ppo_update(self, data): loss_mask data[loss_mask].bool() reward_score data[rewards] # 定义分母布尔掩码 stats_tracker.denominator( n_seqstorch.ones_like(reward_score, dtypetorch.bool), n_valid_tokensloss_mask, ) # 使用分母引用记录张量指标 stats_tracker.stat( advantagesdata[advantages], # [batch, seq_len] kl_rewardsdata[kl_rewards], # [batch, seq_len] denominatorn_valid_tokens ) stats_tracker.stat( task_rewardreward_score.float(), # [batch] seq_lenseqlens.float(), # [batch] denominatorn_seqs )这段代码体现了「先定义分母、再引用分母记录」的核心约定n_seqs按序列粒度统计task_reward、seq_lenn_valid_tokens按 token 粒度统计advantages、kl_rewards。从源码看PPOActor 中还会记录correct_seq_len、incorrect_seq_len、prompt_len、no_eos_ratios、group_loss_weight等序列级指标以及eps_clip、c_clip、use_dual_clip等超参标量全部通过scalar()走默认跟踪器。导出行为训练引擎的export_stats()见 areal/engine/fsdp_engine.py把self.data_parallel_group作为归约组传给export_all()def export_stats(self) - dict[str, float]: # 跨数据并行组 all-reduce return stats_tracker.export_all(reduce_groupself.data_parallel_group) # 所有 DP rank 接收相同的结果底层_aggregate见 stats_tracker.py会按归约类型调用_avg_of、_sum_of、_min_of、_max_of等实现平均值是「掩码内元素求和 ÷ 掩码内元素计数」的加权平均标量则是先 all-reduce 求和、再 all-reduce 计数最终除以计数得到加权均值并同时输出key__count。四、记录 API 与归约类型参考4.1 三种记录方法方法使用场景示例scalar(**kwargs)单个浮点值学习率、超参等scalar(lr0.001, eps0.2)denominator(**kwargs)定义布尔掩码作为分母denominator(validmask.bool())stat(denominator, **kwargs)带掩码的张量指标stat(losstensor, denominatorvalid)源码层面的约束值得注意见 stats_tracker.pydenominator()的值必须是非空的 PyTorch bool 张量否则直接抛ValueErrorstat()的值必须是非空的 float 张量且不允许使用SCALAR归约类型stat()引用的分母必须已经通过denominator()注册过且每次记录的 shape 必须与分母一致内部有assert x.shape y.shape校验所有张量在记录时执行detach().clone()切断与计算图的联系避免梯度残留。4.2 归约类型使用stat()时指标默认归约类型为AVG_MIN_MAX一个键会导出三个输出键stats_tracker.stat(losstensor, denominatorvalid) # 导出{loss/avg: 0.5, loss/min: 0.1, loss/max: 0.9}可用的归约类型由ReduceType枚举定义见 stats_tracker.py类型输出描述AVG_MIN_MAXkey/avg,key/min,key/max张量统计的默认值AVGkey仅加权平均值SUMkey所有元素求和MINkey最小值MAXkey最大值SCALARkey,key__count用于标量值导出加权均值与计数细节上AVG_MIN_MAX的min/max仅在掩码为真的位置取值torch.where(d, v, float(inf)).min()空批次会返回None并从结果中剔除SCALAR类型同时输出key__count正是控制器做跨工作器加权平均所需的计数键。五、作用域、计时与命名跟踪器5.1 层级作用域使用with语句可以嵌套组织指标导出时自动拼接出完整键名with stats_tracker.scope(ppo_actor): with stats_tracker.scope(update): stats_tracker.stat(lossloss_tensor, denominatorvalid) # 键ppo_actor/update/loss/avg实现上Scope上下文管理器在进入/退出时维护一个scope_stack见 stats_tracker.py_get_full_key()把作用域栈与当前键用/拼接。此外还提供了两个配套工具scope_func_wrapper(name)装饰器形式的作用域等价于在函数体内包一层withdisable_scope()临时清空作用域栈用于绕过当前作用域记录全局键。5.2 计时record_timing(key)用time.perf_counter()测量代码块执行时间并固定记录在timeperf/前缀下with stats_tracker.record_timing(rollout): batch actor.prepare_batch(dataloader, workflow) # 键timeperf/rolloutPPOTrainer内部大量使用该机制记录各阶段耗时例如rollout收集批次、critic_values、ref_logp、teacher_logp、rollout_onload、rollout_pause、rollout_offload以及eval等见 areal/trainer/rl_trainer.py 的_onload_rollout/_offload_rollout/train方法。这些timeperf/*指标会以SCALAR类型随训练指标一起导出用于分析训练吞吐与各阶段开销。5.3 命名跟踪器为不同组件隔离指标# 训练指标默认跟踪器 stats_tracker.scalar(grad_norm1.5) # Rollout 指标 stats_tracker.get(rollout).scalar(reward0.8) # 评估指标 stats_tracker.get(eval-rollout).scalar(reward0.9) # 从所有跟踪器导出 all_stats stats_tracker.export_all(reduce_groupgroup)export_all()会通过all_gather_object同步各 rank 已注册的跟踪器名集合从而保证所有 rank 导出的键一致见 stats_tracker.py。另外如果配置了evaluator.eval_before_train: true系统会在第一个训练 step 开始前先执行一次评估并导出指标见 rl_trainer.py用于在微调开始前估计初始模型的性能基线。该配置项定义于 areal/api/cli_args.py。六、完整数据流从收集到记录综合 RLVRWorkflow、RolloutController、FSDPPPOActor 与 PPOTrainer 的实现指标从产生到落盘的完整流程如下Rollout 工作器 训练工作器 ─────────────── ─────────────── workflow.arun_episode() actor.ppo_update(batch) │ │ ▼ ▼ get(rollout).scalar(r0.5) stat(tensor, denommask) │ │ ▼ ▼ export_stats(reduce_groupNone) export_stats(reduce_groupdp_group) {reward: 0.5, reward__count: 1} → all_reduce 跨 DP rank │ │ ▼ │ RolloutController.export_stats() │ → 加权平均跨工作器 │ │ │ └────────────────┬───────────────────────┘ ▼ PPOTrainer._export_and_commit_stats() │ ▼ StatsLogger.commit(stats) │ ┌────────────┼────────────┐ ▼ ▼ ▼ wandb tensorboard swanlabPPOTrainer._export_and_commit_stats()见 rl_trainer.py在每个训练步结束时依次收集三类指标并提交def _export_and_commit_stats(self, epoch, epoch_step, global_step): # 1. 从所有组件收集指标 stats self.actor.export_stats() # 训练指标all-reduced stats.update(self.rollout.export_stats()) # Rollout 指标控制器聚合 stats.update(self.eval_rollout.export_stats()) # 评估指标 # 2. 发送到日志后端仅 rank 0 self.stats_logger.commit(epoch, epoch_step, global_step, stats)值得注意的是eval_rollout的export_stats()同样经由RolloutController的控制器聚合逻辑因此评估指标键前缀eval-rollout/与训练 Rollout 指标键前缀rollout/的聚合路径一致。七、StatsLogger日志后端7.1 职责与生命周期StatsLogger 把聚合后的指标发送到外部日志后端由PPOTrainer在初始化时自动创建StatsLogger(config, ft_spec)见 rl_trainer.py。它在init()中完成各后端的初始化仅在 rank 0 执行以避免重复日志在commit()中过滤内部计数键并写入所有后端在close()中收尾wandb.finish()、swanlab.finish()、trackio.finish()与summary_writer.close()。它还实现了state_dict()/load_state_dict()支持随训练状态一起保存与恢复last_commit_step用于断点续训时保证日志 step 连续。7.2 支持的后端后端配置项描述Weights Biasesconfig.stats_logger.wandb云端实验跟踪SwanLabconfig.stats_logger.swanlab替代实验跟踪TensorBoardconfig.stats_logger.tensorboard本地可视化Trackioconfig.stats_logger.trackioHugging Face 实验跟踪轻量、本地优先说明当前仓库的 StatsLogger 实现 中除文档列举的 wandb / SwanLab / TensorBoard 外还支持 Trackio 后端trackio.init()trackio.log()。Trackio 的配置类TrackioConfig定义于 areal/api/cli_args.py支持disabled/online/local三种模式以及可选的space_id部署 HF Space 远端看板。init()中的细节还包括wandb 支持通过wandb_base_url/wandb_api_key环境注入id_suffixtimestamp时自动追加时间戳以支持多 run 续写resumeallowSwanLab 的api_key缺省时读取SWANLAB_API_KEY环境变量完整实验配置会经redact_sensitive_config()脱敏后连同version_infocommit_id、branch、is_dirty、version一并写入各后端的config字段便于回溯实验环境。7.3 StatsLogger.commit()commit()过滤掉__count内部键后按global_step写入所有后端def commit(self, epoch, step, global_step, data): if dist.is_initialized() and dist.get_rank() ! 0: return # 仅 rank 0 记录 # 过滤掉 __count 键用于内部加权平均 data {k: v for k, v in data.items() if not k.endswith(__count)} # 记录到所有后端 wandb.log(data, stepglobal_step) swanlab.log(data, stepglobal_step) if self.summary_writer: for key, val in data.items(): self.summary_writer.add_scalar(key, val, global_step)同时commit()还会通过tabulate_stats()把指标打印为表格日志print_stats并在日志中输出Epoch/Step/Train step进度信息。_last_commit_step保证断点续训后日志 step 单调递增log_step max(global_step, self._last_commit_step 1)。7.4 配置示例在实验配置中配置日志后端对应StatsLoggerConfig定义于 areal/api/cli_args.pystats_logger: experiment_name: gsm8k_grpo trial_name: run_001 fileroot: /path/to/logs wandb: mode: online # online、offline 或 disabled project: my-project entity: my-team swanlab: mode: online # online、local 或 disabled project: my-project tensorboard: path: /path/to/tensorboard/logs # null 禁用配置字段与源码对应的取值约束experiment_name/trial_name/fileroot为必填项缺失会报错日志根目录由get_log_path()计算为{fileroot}/logs/{user}/{experiment_name}/{trial_name}wandb.mode合法值online、offline、disabled、shared默认disabledwandb.project缺省时回退到experiment_namewandb.name缺省时回退到trial_namegroup缺省为{experiment_name}_{trial_name}swanlab.mode合法值cloud、local、disabled、offline默认disabledproject缺省回退experiment_nametensorboard.path为null时禁用 TensorBoardSummaryWriter不创建trackio.mode合法值disabled、online、local默认disabled。八、最佳实践与常见陷阱选择正确的范式对标量学习率、奖励、超参、耗时使用scalar()对批量 PyTorch 张量通常是训练指标如 loss、advantages、seq_len使用带分母的stat()。先定义分母始终在stat()之前调用denominator()建立掩码关系——stat()引用不存在的分母会直接抛ValueError这是源码层面的硬约束。使用命名跟踪器使用stats_tracker.get(workflow_context.stat_scope()).scalar(...)将 Rolloutrollout和评估eval-rollout指标与训练指标隔离避免键名冲突。注意作用域前缀在with stats_tracker.scope(...)块内记录的所有键都会带上作用域前缀跨组件统计时要注意导出键名的拼接规则。__count键是内部实现细节SCALAR归约会导出key__count用于加权平均commit()会自动过滤控制器合并时也会跳过不要在自定义埋点中刻意手工构造该后缀键。分布式键一致性export_all()通过all_gather_object同步各 rank 的跟踪器集合与键元数据若不同 rank 对同一键的归约类型/分母元数据不一致会在合并时抛出ValueError因此请确保所有 rank 按相同顺序、相同类型记录同一批键。九、测试与验证仓库在 tests/test_stats_tracker.py 中覆盖了指标系统的关键行为包括当某个 rank 缺失某键key_sync_group元数据同步时导出不崩溃test_export_scalar_key_missing_on_this_rank_does_not_crash使用比默认导出组更宽的key_sync_group时默认键仍保持原reduce_group对齐test_export_keeps_default_reduce_group_when_key_sync_group_is_larger跨归约组时 CPU 上的标量会先迁移到 NCCL 设备再执行 all-reducetest_all_reduce_moves_cpu_stat_to_nccl_device_before_reduction。这些测试印证了export()中「per-key reduce_group 覆盖 元数据同步」的设计特定键如 CP-local 场景下的 loss/vocab_* 需要跨 DPCP 归约可以通过stat(..., reduce_group...)单独指定归约组而其余键仍跟随默认组。十、小结AReaL 的指标跟踪系统用「一套 API、两种范式」解决了强化学习训练中最棘手的分布式指标一致性问题流式指标让异步 Rollout 工作器零同步地累积标量导出时由RolloutController按样本数加权聚合批量指标让训练引擎在 DP rank 之间通过 all-reduce 得到全局一致的统计。在此基础上层级作用域、命名跟踪器、六种归约类型、record_timing计时以及StatsLogger对 wandb / SwanLab / TensorBoard / Trackio 的多后端支持构成了从埋点、聚合到可视化的完整闭环。无论是扩展自定义工作流、新增训练指标还是接入新的实验跟踪平台都可以直接从本文的 API 参考与数据流图中获得可落地的指引。【免费下载链接】AReaLThe RL Bridge for LLM-based Agent Applications. Made Simple Flexible.项目地址: https://gitcode.com/GitHub_Trending/are/AReaL创作声明:本文部分内容由AI辅助生成(AIGC),仅供参考