Shortcuts

torchx.runner

Submits AppDef jobs to schedulers.

The runner takes an AppDef (the result of evaluating a component function) along with a scheduler name and run config, and submits it as a job.

from torchx.runner import get_runner

with get_runner() as runner:
    app_handle = runner.run(app, scheduler="kubernetes", cfg=cfg)
    status = runner.status(app_handle)
    print(status)
_images/runner_diagram.png

The Runner submits, monitors, and manages jobs. Use get_runner() to create one with all registered schedulers.

Key methods:

Scheduler instances are created lazily on first use. Use the Runner as a context manager for automatic cleanup.

See Quick Reference for copy-pasteable recipes.

torchx.runner.get_runner(name: str | None = None, component_defaults: dict[str, dict[str, str]] | None = None, **scheduler_params: Any) → Runner[source]

Creates a Runner with all registered schedulers.

with get_runner() as runner:
    app_handle = runner.run(app, scheduler="kubernetes", cfg=cfg)
    print(runner.status(app_handle))
Parameters:

scheduler_params – extra kwargs passed to all scheduler constructors.

class torchx.runner.Runner(name: str = '', scheduler_factories: dict[str, SchedulerFactory] | None = None, component_defaults: dict[str, dict[str, str]] | None = None, scheduler_params: dict[str, object] | None = None)[source]

Submits, monitors, and manages AppDef jobs.

Use get_runner() to create an instance with all registered schedulers.

>>> from torchx.runner import get_runner
>>> runner = get_runner()
>>> runner.scheduler_backends()
['local_cwd', 'local_docker', 'slurm', 'kubernetes', ...]
build_workspace(app_or_role: Role, scheduler: str, cfg: Mapping[str, str | int | float | bool | list[str] | dict[str, str] | None] | None = None) → str | None[source]
build_workspace(app_or_role: AppDef, scheduler: str, cfg: Mapping[str, str | int | float | bool | list[str] | dict[str, str] | None] | None = None) → dict[tuple[str, Workspace], str]

Builds the workspaces and returns the images they were built into.

Does not mutate app_or_role – the build runs against a deep copy. Use this to build once and submit the result many times, instead of letting each run() rebuild:

images = runner.build_workspace(app, "kubernetes", cfg)
for app in per_region_apps:
    pin_workspace_images(app, images)
    runner.run(app, "kubernetes", cfg)

Given a Role, returns that role’s built image, or None if it has no workspace. Given an AppDef, returns {(image, workspace): built} for the roles that have one – keyed by the pair, not by role name, so the result can be pinned onto a different AppDef. A build installs the workspace into role.image, so the base image is half the key: roles sharing a workspace but not an image build to different results, and roles sharing both are built once. Both key halves come from the role as authored, which the build does not overwrite. The pair is the widest key correct for every builder, so one that caches more coarsely than the exact image (on the fbpkg name, say) may rebuild where it could have shared – reuse is missed, never misapplied.

Only workspace is read. The deprecated workspace= argument to run()/dryrun() is applied inside dryrun, too late to prebuild here; set the role attribute instead.

Unlike dryrun(), this neither validates app_or_role nor renders a scheduler request, so it does not fail on request-level problems that are unrelated to building.

A role is omitted (with a warning) when its build wrote to the role beyond image – env, say. Pinning cannot replay those writes, so such a role keeps its workspace and rebuilds on submit. The whole (image, workspace) pair drops out, not just that role: a pair one role has to rebuild is not reusable for the roles that share it.

Returns an empty mapping (or None) when scheduler has no workspace support – there is nothing to build.

cancel(app_handle: str) → None[source]

Requests cancellation. The app transitions to CANCELLED asynchronously.

cfg_from_str(scheduler: str, *cfg_literal: str) → Mapping[str, str | int | float | bool | list[str] | dict[str, str] | None][source]

Convenience function around the scheduler’s runopts.cfg_from_str() method.

Usage:

>>> from torchx.runner import get_runner
>>> runner = get_runner()
>>> runner.cfg_from_str("local_cwd", "log_dir=/tmp/foobar", "prepend_cwd=True")
{'log_dir': '/tmp/foobar', 'prepend_cwd': True}
close() → None[source]

Closes the runner and all scheduler instances. Safe to call multiple times.

delete(app_handle: str) → None[source]

Deletes the app from the scheduler.

describe(app_handle: str) → AppDef | None[source]

Reconstructs the AppDef from the scheduler.

Completeness is scheduler-dependent. Returns None if the app no longer exists.

describe_native(app_handle: str) → AppDryRunInfo | None[source]

Reads the scheduler-native request back for an already-submitted app.

The read-side twin of dryrun(): same AppDryRunInfo wrapper, same scheduler-specific request type. Unlike describe(), the returned request preserves scheduler-specific fields that have no AppDef equivalent. Returns None if the scheduler does not support native read-back or the app no longer exists.

dryrun(app: AppDef, scheduler: str, cfg: Mapping[str, str | int | float | bool | list[str] | dict[str, str] | None] | None = None, workspace: Workspace | str | None = None, parent_run_id: str | None = None) → AppDryRunInfo[source]

Returns what would be submitted without actually submitting.

The returned AppDryRunInfo can be print()-ed for inspection or passed to schedule().

Does not mutate app. The patching this method performs – building each role’s workspace and repointing role.image at the built artifact, injecting the TORCHX_* tracking env vars – lands on a deep copy, which is returned as app. Read the submitted images off that copy, never off the app you passed in:

app = AppDef(roles=[Role(image="foo:latest", workspace=..., ...)])
info = runner.dryrun(app, "kubernetes")

app.roles[0].image       # "foo:latest" -- unpatched, as authored
info.app.roles[0].image  # "foo:<built-workspace-hash>" -- submitted

Role.overrides is the one part not copied: the same dict object is shared by the caller’s role and the copy, so resolving an override through either is visible to both.

Two dryruns of the same app need not render the same request – each rebuilds the workspace, and the build may produce a new image. Re-running from the returned copy does not skip that build either, because workspace is left set after one; it only narrows the difference down to the build, since every other patch is already applied. Clear the workspace to render the exact same request:

info1 = runner.dryrun(app, "kubernetes", cfg)

for role in info1.app.roles:
    role.workspace = None

info2 = runner.dryrun(info1.app, "kubernetes", info1.cfg)
assert info1.request == info2.request
Raises:

ValueError – app has no roles, or a role has no entrypoint or a non-positive num_replicas. Schedulers raise from their own _validate hooks for backend-specific violations.

dryrun_component(component: str, component_args: list[str] | dict[str, Any], scheduler: str, cfg: Mapping[str, str | int | float | bool | list[str] | dict[str, str] | None] | None = None, workspace: Workspace | str | None = None, parent_run_id: str | None = None) → AppDryRunInfo[source]

Like run_component() but returns the request without submitting.

list(scheduler: str, cfg: Mapping[str, str | int | float | bool | list[str] | dict[str, str] | None] | None = None) → list[ListAppResponse][source]

Lists jobs on the scheduler.

Parameters:

cfg – scheduler config, used by some schedulers for backend routing.

log_lines(app_handle: str, role_name: str, k: int = 0, regex: str | None = None, since: datetime | None = None, until: datetime | None = None, should_tail: bool = False, streams: Stream | None = None) → Iterable[str][source]

Returns an iterator over log lines for the k-th replica of a role.

Important

k is the node (host) id, NOT the worker rank.

Warning

Completeness is scheduler-dependent. Lines may be partial or missing if logs have been purged. Do not use this for programmatic output parsing.

Lines include trailing whitespace (\n). Use print(line, end="") to avoid double newlines.

Parameters:
  • k – replica (node) index

  • regex – optional filter pattern

  • since – start cursor (scheduler-dependent)

  • until – end cursor (scheduler-dependent)

run(app: AppDef, scheduler: str, cfg: Mapping[str, str | int | float | bool | list[str] | dict[str, str] | None] | None = None, workspace: Workspace | str | None = None, parent_run_id: str | None = None, *, dryrun: Literal[True]) → AppDryRunInfo[source]
run(app: AppDef, scheduler: str, cfg: Mapping[str, str | int | float | bool | list[str] | dict[str, str] | None] | None = None, workspace: Workspace | str | None = None, parent_run_id: str | None = None, *, dryrun: Literal[False] = False) → str

Submits an AppDef or returns its dry-run info.

# Submit
handle = runner.run(app, "mkube", cfg=cfg)

# Dryrun — inspect without submitting
info = runner.run(app, "mkube", cfg=cfg, dryrun=True)
print(info)

Does not mutate app – see dryrun(), which this delegates to. The workspace build and env injection land on an internal deep copy: with dryrun=True that copy comes back as info.app; with dryrun=False it is submitted and dropped, and only the AppHandle is returned.

Parameters:

dryrun – If True, only validate and render the request without submitting. Returns AppDryRunInfo. If False (default), submit and return the AppHandle.

Raises:

ValueError – propagated from dryrun().

run_component(component: str, component_args: list[str] | dict[str, Any], scheduler: str, cfg: Mapping[str, str | int | float | bool | list[str] | dict[str, str] | None] | None = None, workspace: Workspace | str | None = None, parent_run_id: str | None = None) → str[source]

Resolves and runs a named component.

component resolution order (high → low):

  1. User-registered torchx.components entry points

  2. Builtins relative to torchx.components (e.g. "dist.torchrun")

  3. File-based path/to/file.py:function_name

schedule(dryrun_info: AppDryRunInfo) → str[source]

Submits a previously dry-run request, allowing request mutation.

dryrun_info = runner.dryrun(app, scheduler="kubernetes", cfg)
dryrun_info.request.foo = "bar"  # mutate the raw request
app_handle = runner.schedule(dryrun_info)

Warning

Use sparingly. Overwriting many raw scheduler fields may cause your usage to diverge from TorchX’s supported API.

Only dryrun_info.request is submitted; edits to dryrun_info.app made after dryrun() returned are not re-rendered into it.

scheduler_backends() → list[str][source]

Returns all registered scheduler backend names.

scheduler_run_opts(scheduler: str) → runopts[source]

Returns the runopts for the given scheduler.

status(app_handle: str) → AppStatus | None[source]

Returns app status, or None if the app no longer exists.

stop(app_handle: str) → None[source]

Deprecated since version Use: cancel() instead.

wait(app_handle: str, wait_interval: float = 10) → AppStatus | None[source]

Blocks until the app reaches a terminal state.

Parameters:

wait_interval – seconds between status polls

See also

Quick Reference

Single-page reference with imports, types, and copy-pasteable recipes.

torchx.schedulers

Scheduler API reference and implementation guide.

.torchxconfig

Configuring scheduler options via .torchxconfig.

Docs

Access comprehensive developer documentation for PyTorch

View Docs

Tutorials

Get in-depth tutorials for beginners and advanced developers

View Tutorials

Resources

Find development resources and get your questions answered

View Resources