Summary
art.tokenize parallelizes through src/art/trajectories/_parallel.py, but on a 64-CPU driver the steady-state ceiling is ~84k tokens/s because the process pool is hard-capped at _PROCESS_MAX_WORKERS = 4 and the thread path is GIL-bound at ~30k tokens/s. Proposing caladan experiment 059's per-step tokenization as a concrete benchmark workload for improving this.
Benchmark workload (caladan experiments/059-robust-policy.py)
Per training step the driver tokenizes two things back to back:
art.tokenize(groups, multi_history=True, base_model="Qwen/Qwen3.6-35B-A3B") — 8 TrajectoryGroups x 5 trajectories (1 real corpus trajectory + 4 dynamics-model rollouts), each rollout carrying two ~20k-token histories (policy view + adversary/dynamics view with captured sampled tokens and logprobs). ~40 leaves, ~1.5M tokens.
art.tokenize([cal.dynamics.trajectory(t, model=adversary) for t in reals], model=adversary, base_model=...) — 8 single-history adversary views of the real trajectories, ~20k tokens each.
The corpus is tau-bench retail-49 transcripts (corpus4096, ~7 MB per trajectory materialized). A self-contained reproduction only needs 8–16 of those trajectories; happy to share a slice.
Measurements (this box: 120 CPUs, cgroup limit 64, tokenizer cached)
16 adversary-view trajectories, ~300k tokens, art.tokenize(..., model="adv", base_model="Qwen/Qwen3.6-35B-A3B"):
| run |
wall |
throughput |
backend |
| 0 |
11.8 s |
26k tok/s |
threads (4 workers) |
| 1 |
9.2 s |
33k tok/s |
threads |
| 2 |
26.7 s |
11k tok/s |
process pool spawn + tokenizer load in 4 workers |
| 3 |
3.6 s |
84k tok/s |
processes (4 workers) |
| 8 + 8 sequential, warm |
4.0 s |
|
processes |
8 + 8 asyncio.gather, warm |
3.5 s |
|
processes, shared pool |
Observations from _parallel.py
_PROCESS_MAX_WORKERS = 4 regardless of _cpu_capacity(); with 64 CPUs available, 60 sit idle.
- Processes are only enabled after a tuning key has >= 2 thread runs each >= 1 s (
_consider_processes). The tuning key includes _size_bucket(len(values)), so every distinct batch-size bucket pays that warm-up again, and the first process run also pays a ~25 s spawn + tokenizer load.
- Thread path:
trajectory.tokenize is mostly Python (chat template render, message/sampled-token matching, flag construction); only the Rust encode releases the GIL, so workers > ~4 doesn't help.
- Process path: one pickle round trip per trajectory (
_process_payloads serial on one thread, _deserialize_process_result + _rebind_process_result per item). For ~7 MB trajectories this is a nontrivial share of the 3.6 s.
- The autotuner (
_observe_process) assumes exclusive use of the pool; concurrent art.tokenize calls contend for the same 4 workers and skew each other's rate measurements (it disables processes if they don't beat threads by 5%).
Ideas
- Derive the process cap from capacity (e.g.
min(capacity, 16)), or expose it (ART_TOKENIZE_PROCESSES).
- Let callers opt into processes up front instead of waiting for two slow thread runs per size bucket; or warm the pool at import/first call in the background.
- Batch several trajectories per process task to amortize pickling and the per-task scheduling overhead.
- Reuse the pickled payload for
tensorize rather than re-pickling.
- Make the tuner concurrency-aware (share
next_workers across in-flight calls, or measure per-pool rather than per-call).
Would be glad to run any candidate against the 059 shapes above.
Summary
art.tokenizeparallelizes throughsrc/art/trajectories/_parallel.py, but on a 64-CPU driver the steady-state ceiling is ~84k tokens/s because the process pool is hard-capped at_PROCESS_MAX_WORKERS = 4and the thread path is GIL-bound at ~30k tokens/s. Proposing caladan experiment 059's per-step tokenization as a concrete benchmark workload for improving this.Benchmark workload (caladan
experiments/059-robust-policy.py)Per training step the driver tokenizes two things back to back:
art.tokenize(groups, multi_history=True, base_model="Qwen/Qwen3.6-35B-A3B")— 8TrajectoryGroups x 5 trajectories (1 real corpus trajectory + 4 dynamics-model rollouts), each rollout carrying two ~20k-token histories (policy view + adversary/dynamics view with captured sampled tokens and logprobs). ~40 leaves, ~1.5M tokens.art.tokenize([cal.dynamics.trajectory(t, model=adversary) for t in reals], model=adversary, base_model=...)— 8 single-history adversary views of the real trajectories, ~20k tokens each.The corpus is tau-bench retail-49 transcripts (
corpus4096, ~7 MB per trajectory materialized). A self-contained reproduction only needs 8–16 of those trajectories; happy to share a slice.Measurements (this box: 120 CPUs, cgroup limit 64, tokenizer cached)
16 adversary-view trajectories, ~300k tokens,
art.tokenize(..., model="adv", base_model="Qwen/Qwen3.6-35B-A3B"):asyncio.gather, warmObservations from
_parallel.py_PROCESS_MAX_WORKERS = 4regardless of_cpu_capacity(); with 64 CPUs available, 60 sit idle._consider_processes). The tuning key includes_size_bucket(len(values)), so every distinct batch-size bucket pays that warm-up again, and the first process run also pays a ~25 s spawn + tokenizer load.trajectory.tokenizeis mostly Python (chat template render, message/sampled-token matching, flag construction); only the Rust encode releases the GIL, soworkers > ~4doesn't help._process_payloadsserial on one thread,_deserialize_process_result+_rebind_process_resultper item). For ~7 MB trajectories this is a nontrivial share of the 3.6 s._observe_process) assumes exclusive use of the pool; concurrentart.tokenizecalls contend for the same 4 workers and skew each other's rate measurements (it disables processes if they don't beat threads by 5%).Ideas
min(capacity, 16)), or expose it (ART_TOKENIZE_PROCESSES).tensorizerather than re-pickling.next_workersacross in-flight calls, or measure per-pool rather than per-call).Would be glad to run any candidate against the 059 shapes above.