Skip to content

Commit 7af08d2

Browse files
ko3n1gclaude
andcommitted
fix(kubeflow): stream logs once, not per replica
torchx calls scheduler.log_iter(app_id, role_name, k=...) once per replica (k = 0..num_nodes-1). The Kubeflow log_iter ignored k and re-ran fetch_logs — which tails the entire jobset via the jobset-name selector — for every replica, producing N independent tail streams (each with its own dedup state) and N-fold-duplicating every console line (prefixed <role>/<k>). At 16 nodes that's 16x the log volume, which also overruns the CI job-log limit on long runs. Stream only for k == 0; that single tail already covers all ranks (and writes log-allranks_0.out once). Co-Authored-By: Claude Opus 4.8 (1M context) <noreply@anthropic.com> Signed-off-by: oliver könig <okoenig@nvidia.com>
1 parent 636ec99 commit 7af08d2

1 file changed

Lines changed: 8 additions & 0 deletions

File tree

‎nemo_run/run/torchx_backend/schedulers/kubeflow.py‎

Lines changed: 8 additions & 0 deletions
Original file line numberDiff line numberDiff line change
@@ -188,6 +188,14 @@ def log_iter(
188188
if not executor:
189189
return []
190190

191+
# fetch_logs tails ALL pods of the jobset in a single call (it powers the
192+
# log-allranks_0.out capture and the cross-rank dedup). torchx invokes
193+
# log_iter once per replica (k = 0..num_nodes-1); streaming on every k
194+
# would re-tail the whole jobset N times — each tail with its own dedup
195+
# state — and N×-duplicate every console line. Stream only for k == 0.
196+
if k != 0:
197+
return []
198+
191199
logs = executor.fetch_logs(job_name=job_name, stream=should_tail)
192200
if isinstance(logs, str):
193201
if len(logs) == 0:

0 commit comments

Comments
 (0)