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)
The Runner submits, monitors, and manages jobs. Use
get_runner() to create one with all registered schedulers.
Key methods:
run()/run_component()– submit a jobstatus()– poll current statewait()– block until terminal statecancel()– request cancellationdelete()– remove a job definition from the schedulerlog_lines()– stream log outputlist()– list jobs on a schedulerdryrun()– preview what would be submitted without submittingschedule()– submit a previously dry-run request (allows request mutation)
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
Runnerwith 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
AppDefjobs.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, orNoneif it has noworkspace. Given anAppDef, 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 differentAppDef. A build installs the workspace intorole.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
workspaceis read. The deprecatedworkspace=argument torun()/dryrun()is applied insidedryrun, 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
CANCELLEDasynchronously.
- 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}
- describe(app_handle: str) AppDef | None[source]¶
Reconstructs the
AppDeffrom the scheduler.Completeness is scheduler-dependent. Returns
Noneif 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(): sameAppDryRunInfowrapper, same scheduler-specificrequesttype. Unlikedescribe(), the returnedrequestpreserves scheduler-specific fields that have noAppDefequivalent. ReturnsNoneif 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
AppDryRunInfocan beprint()-ed for inspection or passed toschedule().Does not mutate app. The patching this method performs – building each role’s
workspaceand repointingrole.imageat the built artifact, injecting theTORCHX_*tracking env vars – lands on a deep copy, which is returned asapp. 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.overridesis 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
workspaceis 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
entrypointor a non-positivenum_replicas. Schedulers raise from their own_validatehooks 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
kis 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). Useprint(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
AppDefor 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: withdryrun=Truethat copy comes back asinfo.app; withdryrun=Falseit is submitted and dropped, and only theAppHandleis returned.- Parameters:
dryrun – If
True, only validate and render the request without submitting. ReturnsAppDryRunInfo. IfFalse(default), submit and return theAppHandle.- 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.
componentresolution order (high → low):User-registered
torchx.componentsentry pointsBuiltins relative to
torchx.components(e.g."dist.torchrun")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.requestis submitted; edits todryrun_info.appmade afterdryrun()returned are not re-rendered into it.
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.