Skip to content
Merged
Show file tree
Hide file tree
Changes from all commits
Commits
File filter

Filter by extension

Filter by extension


Conversations
Failed to load comments.
Loading
Jump to
Jump to file
Failed to load files.
Loading
Diff view
Diff view
122 changes: 72 additions & 50 deletions docs/sphinx/source/en/2-user_guide/1-training/3-logging.md

Large diffs are not rendered by default.

4 changes: 2 additions & 2 deletions docs/sphinx/source/en/2-user_guide/2-algorithms/2-appo.md
Original file line number Diff line number Diff line change
Expand Up @@ -60,10 +60,10 @@ gantt
Weight Publish → collector :crit, l4, 28000, 30000

section Iter Wall
perf/iter_ms (learner loop only) : l5, 12000, 30000
Perf/iteration_time (learner loop only) : l5, 12000, 30000
```

> The axis is schematic (relative, not real-ms). The collector subprocess produces rollouts through the 4-slot ring buffer in parallel with the learner, so **Collector Wait ≈ 0** in steady state. `perf/iter_ms` counts only this learner loop (it includes Collector Wait but not the collector's parallel rollout compute); the red Weight Publish marks the end of the iteration when fresh weights are published to the collector. Field meanings are on the [logging page](../1-training/3-logging.md).
> The axis is schematic (relative, not real-ms). The collector subprocess produces rollouts through the 4-slot ring buffer in parallel with the learner, so **Collector Wait ≈ 0** in steady state. `Perf/iteration_time` counts only this learner loop (it includes Collector Wait but not the collector's parallel rollout compute); the red Weight Publish marks the end of the iteration when fresh weights are published to the collector. Field meanings are on the [logging page](../1-training/3-logging.md).

## Key Fields

Expand Down
118 changes: 68 additions & 50 deletions docs/sphinx/source/zh_CN/2-user_guide/1-training/3-logging.md

Large diffs are not rendered by default.

Original file line number Diff line number Diff line change
Expand Up @@ -58,10 +58,10 @@ gantt
Weight Publish 写回 collector :crit, l4, 28000, 30000

section Iter Wall
perf/iter_ms(仅 learner 这圈) : l5, 12000, 30000
Perf/iteration_time(仅 learner 这圈) : l5, 12000, 30000
```

> 横轴为示意相对时长(非真实 ms 比例)。collector 子进程经 4 槽 ring buffer 与 learner 并行产出 rollout,稳态下 **Collector Wait ≈ 0**。`perf/iter_ms` 仅计 learner 这一圈(含 Collector Wait,但不含 collector 的并行采集计算);红色 Weight Publish 标志该轮迭代结束、向 collector 发布新权重。各字段含义见[日志页](../1-training/3-logging.md)。
> 横轴为示意相对时长(非真实 ms 比例)。collector 子进程经 4 槽 ring buffer 与 learner 并行产出 rollout,稳态下 **Collector Wait ≈ 0**。`Perf/iteration_time` 仅计 learner 这一圈(含 Collector Wait,但不含 collector 的并行采集计算);红色 Weight Publish 标志该轮迭代结束、向 collector 发布新权重。各字段含义见[日志页](../1-training/3-logging.md)。

## 关键字段

Expand Down
6 changes: 3 additions & 3 deletions docs/sphinx/source/zh_CN/2-user_guide/2-algorithms/3-sac.md
Original file line number Diff line number Diff line change
Expand Up @@ -66,9 +66,9 @@ uv run train --algo sac --task g1_walk_flat --sim mujoco \

rank 0 独占终端与 TensorBoard/W&B logger;其他 learner 不刷新终端,也不创建独立
tfevents。rank 0 在每个 learner iteration 汇总跨 rank 标量:loss、reward 和 timing
取均值,计数与并行吞吐求和。`perf/steps_per_sec` / 终端 `Steps/s` 表示所有
collector 的聚合 env-step 吞吐,`perf/effective_samples_per_sec` / 终端 `Samples/s`
表示所有 learner 的聚合有效样本吞吐。checkpoint 同样由 rank 0 独占:每个保存间隔和
取均值,计数与并行吞吐求和。`Perf/total_fps` / 终端 `Steps/s` 表示所有
collector 的聚合 env-step 吞吐;终端 `Rows/s` 表示所有 learner 的聚合 replay 行
吞吐(仅终端显示,不持久化;可由 run 配置与 `Perf/iteration_time` 推导)。checkpoint 同样由 rank 0 独占:每个保存间隔和
训练结束只在 canonical run 目录写一份模型;其他 rank 复用该路径完成进程协调,不创建
rank 子目录或任何日志文件。
自动生成的多卡 run 目录以 `_gpuxN` 结尾(例如 `_gpux2`);单卡目录和显式
Expand Down
4 changes: 2 additions & 2 deletions pyproject.rocm.toml
Original file line number Diff line number Diff line change
Expand Up @@ -72,12 +72,12 @@ mujoco = [
"mjbatch-uni~=0.2.1",
]
motrix = ["unisim-core[motrix]>=1.7.8"]
uni_rl = ["unilab-rl==1.4.0"]
uni_rl = ["unilab-rl==1.4.1"]

[dependency-groups]
# Dev/CI exercises the APPO/off-policy/multi-GPU paths, so the optional
# unilab-rl runtime stays installed in the dev environment.
dev = ["pytest", "pytest-cov", "ruff", "mypy", "pyright>=1.1.408", "unilab-rl==1.4.0"]
dev = ["pytest", "pytest-cov", "ruff", "mypy", "pyright>=1.1.408", "unilab-rl==1.4.1"]

[[tool.uv.index]]
name = "pytorch-rocm72"
Expand Down
4 changes: 2 additions & 2 deletions pyproject.toml
Original file line number Diff line number Diff line change
Expand Up @@ -109,7 +109,7 @@ newton = [
]
motrix = ["unisim-core[motrix]>=1.7.8"]
genesis = ["genesis-world==1.3.3"]
uni_rl = ["unilab-rl==1.4.0"]
uni_rl = ["unilab-rl==1.4.1"]
# Wheels are CPython 3.12/3.13 Linux x86_64 only; the marker keeps uv lock
# resolvable for the other required-environments.
superdex = [
Expand All @@ -119,7 +119,7 @@ superdex = [
[dependency-groups]
# Dev/CI exercises the APPO/off-policy/multi-GPU paths, so the optional
# unilab-rl runtime stays installed in the dev environment.
dev = ["pytest", "pytest-cov", "ruff", "mypy", "pyright>=1.1.408", "unilab-rl==1.4.0"]
dev = ["pytest", "pytest-cov", "ruff", "mypy", "pyright>=1.1.408", "unilab-rl==1.4.1"]

[[tool.uv.index]]
name = "r2-cu130"
Expand Down
104 changes: 82 additions & 22 deletions scripts/benchmark/rl/benchmark_offpolicy_dp_scaling.py
Original file line number Diff line number Diff line change
Expand Up @@ -12,17 +12,22 @@
Measurement conventions:

- Collector throughput (Steps/s): mean of the last 50% of canonical rank-0
``perf/steps_per_sec`` samples. In DP runs the runner already sums the
per-rank collector rates before logging.
- Learner throughput (Samples/s): the same tail mean over
``perf/effective_samples_per_sec``. This counts effective learner samples
(including configured sample multipliers), summed across ranks.
``Perf/total_fps`` samples (legacy ``perf/steps_per_sec`` for pre-1.4.1
event files). In DP runs the runner already sums the per-rank collector
rates before logging.
- Learner throughput (Samples/s): the same tail mean over the legacy
``perf/effective_samples_per_sec`` tag. unilab-rl 1.4.1 no longer persists
learner replay throughput, so runs without that tag derive it per iteration
as ``world_size * algo.batch_size * algo.updates_per_step /
Perf/iteration_time`` (run configuration from ``run_config.json``). This
counts effective learner replay rows (including configured sample
multipliers), summed across ranks.
- Both fields report their own N-way / N=1 scaling ratio. The roadmap verdict
remains attached to collector Steps/s: ``pass`` at >= 1.7
(``SCALING_PASS_THRESHOLD``), otherwise ``below threshold``. The verdict is
data only and does not affect the exit code.
- The tail uses ceil(n/2) points (with n samples, ``n - n//2``), skipping
collector warm-up and replay prefill. Missing either throughput tag is a
collector warm-up and replay prefill. Missing either throughput series is a
hard error, never silently skipped.
- Exit code is non-zero only when a run itself failed (subprocess error,
non-completed summary, or missing artifacts).
Expand Down Expand Up @@ -65,10 +70,23 @@
# --sim mujoco` (see src/unilab/cli.py build_route for off-policy algos).
ROUTE_OVERRIDES = ("task=g1_walk_flat/mujoco",)

STEPS_PER_SEC_TAG = "perf/steps_per_sec"
SAMPLES_PER_SEC_TAG = "perf/effective_samples_per_sec"
REWARD_TAG = "reward/mean"
DP_SYNC_TIME_TAG = "train/dp_sync_time"
# Scalar tag candidates are ``(tag, scale)`` pairs ordered newest schema
# first. unilab-rl 1.4.1 introduced the canonical metric schema
# (``uni_rl.logging.metric_schema``; migration table in unilab_rl
# ``docs/metrics.md``) without rewriting historical event files, so the first
# tag present in a run wins and its values are scaled into the units the
# report columns already use (per-second rates, seconds for durations).
STEPS_PER_SEC_TAGS = (("Perf/total_fps", 1.0), ("perf/steps_per_sec", 1.0))
REWARD_TAGS = (("Train/mean_reward", 1.0), ("reward/mean", 1.0))
DP_SYNC_TIME_TAGS = (("Perf/dp_gradient_sync_ms_per_rank", 0.001), ("train/dp_sync_time", 1.0))
ITERATION_TIME_TAGS = (("Perf/iteration_time", 1.0), ("perf/iter_ms", 0.001))
# unilab-rl 1.4.1 removed the persisted learner replay throughput chart with
# no replacement tag; post-1.4.1 runs derive it from the run configuration
# and ``Perf/iteration_time`` (see derive_learner_samples_per_sec).
SAMPLES_PER_SEC_TAGS = (
("perf/effective_samples_per_sec", 1.0),
("perf/learner_samples_per_sec", 1.0),
)

DEFAULT_ITERATIONS = 300
DEFAULT_DEVICES = "0,1"
Expand Down Expand Up @@ -117,8 +135,8 @@ def find_event_files(rank_dir: Path) -> list[Path]:
return sorted(rank_dir.glob("events.out.tfevents.*"))


def read_scalar_series(rank_dir: Path, tag: str) -> list[float]:
"""Scalar values of ``tag`` from a rank's tfevents (in event order).
def read_scalar_series(rank_dir: Path, candidates: Sequence[tuple[str, float]]) -> list[float]:
"""Scaled values of the first present tag among ``candidates``.

An absent tag yields ``[]``; a rank directory without any tfevents file
raises ``RunParseError`` so missing rank data is never silent.
Expand All @@ -134,9 +152,37 @@ def read_scalar_series(rank_dir: Path, tag: str) -> list[float]:

accumulator = event_accumulator.EventAccumulator(str(event_files[0]))
accumulator.Reload()
if tag not in accumulator.Tags()["scalars"]:
tags = accumulator.Tags()["scalars"]
for tag, scale in candidates:
if tag in tags:
return [float(event.value) * scale for event in accumulator.Scalars(tag)]
return []


def derive_learner_samples_per_sec(run_dir: Path, world_size: int) -> list[float]:
"""Reconstruct aggregate learner replay rows/s for post-1.4.1 runs.

unilab-rl 1.4.1 no longer persists learner replay throughput; per
unilab_rl ``docs/metrics.md`` it derives from the run configuration —
``algo.batch_size * algo.updates_per_step`` replay rows per rank per
iteration, summed across ranks — divided by the measured
``Perf/iteration_time``.
"""
iteration_times = read_scalar_series(run_dir, ITERATION_TIME_TAGS)
config_path = Path(run_dir) / "run_config.json"
if not iteration_times or not config_path.is_file():
return []
config = json.loads(config_path.read_text(encoding="utf-8")).get("config") or {}
algo = config.get("algo") or {}
try:
rows_per_rank_per_iter = int(algo["batch_size"]) * int(algo["updates_per_step"])
except (KeyError, TypeError, ValueError):
return []
return [float(event.value) for event in accumulator.Scalars(tag)]
return [
world_size * rows_per_rank_per_iter / iteration_time
for iteration_time in iteration_times
if iteration_time > 0
]


def parse_run(run_dir: Path, world_size: int) -> dict[str, Any]:
Expand All @@ -160,14 +206,24 @@ def parse_run(run_dir: Path, world_size: int) -> dict[str, Any]:
f"run {run_dir} did not complete: status={status!r} error={summary.get('error')!r}"
)

collector_series = read_scalar_series(run_dir, STEPS_PER_SEC_TAG)
collector_series = read_scalar_series(run_dir, STEPS_PER_SEC_TAGS)
if not collector_series:
raise RunParseError(f"run has no {STEPS_PER_SEC_TAG!r} samples under {run_dir}")
learner_series = read_scalar_series(run_dir, SAMPLES_PER_SEC_TAG)
raise RunParseError(
f"run has no collector throughput samples "
f"({' or '.join(tag for tag, _ in STEPS_PER_SEC_TAGS)}) under {run_dir}"
)
learner_series = read_scalar_series(run_dir, SAMPLES_PER_SEC_TAGS)
learner_throughput_source = "tfevents"
if not learner_series:
raise RunParseError(f"run has no {SAMPLES_PER_SEC_TAG!r} samples under {run_dir}")
dp_sync_samples = read_scalar_series(run_dir, DP_SYNC_TIME_TAG)
reward_series = read_scalar_series(run_dir, REWARD_TAG)
learner_series = derive_learner_samples_per_sec(run_dir, world_size)
learner_throughput_source = "derived"
if not learner_series:
raise RunParseError(
f"run has no 'perf/effective_samples_per_sec' samples and learner "
f"throughput cannot be derived under {run_dir}"
)
dp_sync_samples = read_scalar_series(run_dir, DP_SYNC_TIME_TAGS)
reward_series = read_scalar_series(run_dir, REWARD_TAGS)
return {
"run_dir": str(run_dir),
"world_size": world_size,
Expand All @@ -176,6 +232,7 @@ def parse_run(run_dir: Path, world_size: int) -> dict[str, Any]:
"training_wall_time_sec": summary.get("training_wall_time_sec"),
"num_collector_throughput_samples": len(collector_series),
"num_learner_throughput_samples": len(learner_series),
"learner_throughput_source": learner_throughput_source,
"steady_state_collector_steps_per_s": steady_state_mean(collector_series),
"steady_state_learner_samples_per_s": steady_state_mean(learner_series),
"final_mean_reward": reward_series[-1] if reward_series else None,
Expand Down Expand Up @@ -468,8 +525,11 @@ def main(argv: list[str] | None = None) -> int:
"scaling_pass_threshold": SCALING_PASS_THRESHOLD,
"steady_state_tail_fraction": STEADY_STATE_TAIL_FRACTION,
"throughput_tags": {
"collector_steps_per_sec": STEPS_PER_SEC_TAG,
"learner_samples_per_sec": SAMPLES_PER_SEC_TAG,
"collector_steps_per_sec": [tag for tag, _ in STEPS_PER_SEC_TAGS],
"learner_samples_per_sec": [tag for tag, _ in SAMPLES_PER_SEC_TAGS],
"learner_samples_per_sec_derived": (
"world_size * algo.batch_size * algo.updates_per_step / Perf/iteration_time"
),
},
},
},
Expand Down
91 changes: 65 additions & 26 deletions scripts/benchmark/rl/extract_offpolicy_metrics.py
Original file line number Diff line number Diff line change
@@ -1,5 +1,18 @@
#!/usr/bin/env python3
"""Extract end-to-end off-policy benchmark metrics from TensorBoard event files."""
"""Extract end-to-end off-policy benchmark metrics from TensorBoard event files.

unilab-rl 1.4.1 replaced the ``timing/*`` and lowercase ``perf/*`` tags with the
canonical schema in ``uni_rl.logging.metric_schema`` (see ``docs/metrics.md`` in
the unilab_rl repo). Historical event files are immutable, so every row below
lists ``(tag, scale)`` candidates newest-first: the first tag present in the
file wins and its values are multiplied by ``scale`` to keep each row's
long-standing unit (milliseconds / per-second rates). Two canonical fields are
seconds where the retired tags were milliseconds: ``Perf/iteration_time``
(was ``perf/iter_ms``) and ``Perf/learning_time`` (was
``timing/learner_train_ms``). Rows whose source chart was retired without a
replacement (collector active throughput, collector cycle total) only carry
their legacy tag and report NaN on post-1.4.1 runs.
"""

from __future__ import annotations

Expand All @@ -9,20 +22,44 @@

from tensorboard.backend.event_processing import event_accumulator

TAGS = {
"iter_ms": "perf/iter_ms",
"steps_per_sec": "perf/steps_per_sec",
"collector_active_steps_per_sec": "perf/collector_active_steps_per_sec",
"collector_cycle_ms": "perf/collector_cycle_ms",
"collector_learner_action_wait_ms": "timing/collector_learner_action_wait_ms",
"collector_replay_write_ms": "timing/collector_replay_write_ms",
"learner_collector_wait_ms": "timing/learner_collector_wait_ms",
"learner_inference_ms": "timing/learner_inference_ms",
"learner_replay_batch_wait_ms": "timing/learner_replay_batch_wait_ms",
"learner_replay_sample_ms": "timing/learner_replay_sample_ms",
"replay_ingress_h2d_submit_ms": "timing/replay_ingress_h2d_submit_ms",
"learner_collector_release_ms": "timing/learner_collector_release_ms",
"learner_train_ms": "timing/learner_train_ms",
TAGS: dict[str, tuple[tuple[str, float], ...]] = {
"iter_ms": (("Perf/iteration_time", 1000.0), ("perf/iter_ms", 1.0)),
"steps_per_sec": (("Perf/total_fps", 1.0), ("perf/steps_per_sec", 1.0)),
"collector_active_steps_per_sec": (("perf/collector_active_steps_per_sec", 1.0),),
"collector_cycle_ms": (("perf/collector_cycle_ms", 1.0),),
"collector_learner_action_wait_ms": (
("Perf/collector_learner_action_wait_ms", 1.0),
("timing/collector_learner_action_wait_ms", 1.0),
),
"collector_replay_write_ms": (
("Perf/collector_replay_write_ms", 1.0),
("timing/collector_replay_write_ms", 1.0),
),
"learner_collector_wait_ms": (
("Perf/learner_collector_wait_ms", 1.0),
("timing/learner_collector_wait_ms", 1.0),
),
"learner_inference_ms": (
("Perf/learner_inference_ms", 1.0),
("timing/learner_inference_ms", 1.0),
),
"learner_replay_batch_wait_ms": (
("Perf/learner_replay_batch_wait_ms", 1.0),
("timing/learner_replay_batch_wait_ms", 1.0),
),
"learner_replay_sample_ms": (
("Perf/learner_replay_sample_ms", 1.0),
("timing/learner_replay_sample_ms", 1.0),
),
"replay_ingress_h2d_submit_ms": (
("Perf/replay_ingress_h2d_submit_ms", 1.0),
("timing/replay_ingress_h2d_submit_ms", 1.0),
),
"learner_collector_release_ms": (
("Perf/learner_collector_release_ms", 1.0),
("timing/learner_collector_release_ms", 1.0),
),
"learner_train_ms": (("Perf/learning_time", 1000.0), ("timing/learner_train_ms", 1.0)),
}


Expand All @@ -41,6 +78,18 @@ def _average_last(scalars: list[event_accumulator.ScalarEvent], n: int) -> float
return sum(values[-n:]) / len(values[-n:])


def _extract_row(
ea: event_accumulator.EventAccumulator,
candidates: tuple[tuple[str, float], ...],
last: int,
) -> float:
available = set(ea.Tags()["scalars"])
for tag, scale in candidates:
if tag in available:
return _average_last(ea.Scalars(tag), last) * scale
return float("nan")


def main(argv: list[str]) -> int:
parser = argparse.ArgumentParser(description=__doc__)
parser.add_argument("log_dir", type=Path)
Expand All @@ -56,17 +105,7 @@ def main(argv: list[str]) -> int:
ea = event_accumulator.EventAccumulator(str(event_file))
ea.Reload()

available_tags = {tag: False for tag in TAGS.values()}
for tag in ea.Tags()["scalars"]:
available_tags[tag] = True

rows = []
for label, tag in TAGS.items():
if not available_tags.get(tag, False):
rows.append((label, float("nan")))
continue
scalars = ea.Scalars(tag)
rows.append((label, _average_last(scalars, args.last)))
rows = [(label, _extract_row(ea, candidates, args.last)) for label, candidates in TAGS.items()]

if args.json:
import json
Expand Down
Loading
Loading