Skip to content

Pipelines

View Markdown

A Pipeline runs Jobs as stages. A stage starts when the stages it depends on have succeeded, reads their outputs as inputs, and runs in parallel with every other stage that is ready. The whole Pipeline shares one spending cap.

This Pipeline prepares a dataset on a CPU, trains on a GPU, and evaluates the trained model:

pipeline.yaml
apiVersion: nodus.dev/v1
kind: Pipeline
metadata:
name: train-eval
spec:
# One cap for the whole run: every stage's spend counts against it.
maxCostUSD: "0.50"
failurePolicy: FailFast
# Stage templates inherit these fields unless they set their own.
defaults:
image: nodus/pytorch
timeout: 15m
stages:
- name: prepare
template:
resources:
cpu: "2"
memory: 4Gi
command:
- python
- -c
- |
import os, torch
os.makedirs("/nodus/outputs/shards", exist_ok=True)
x = torch.randn(4096, 16)
y = (x.sum(dim=1, keepdim=True) > 0).float()
torch.save({"x": x, "y": y}, "/nodus/outputs/shards/data.pt")
outputs:
- name: shards
path: /nodus/outputs/shards
- name: train
dependsOn: [prepare]
inputs:
# Mounted at /nodus/inputs/data and named by NODUS_INPUT_DATA.
- name: data
fromStage: prepare
output: shards
template:
resources:
gpu: L4
command:
- python
- -c
- |
import os, torch
d = torch.load(os.path.join(os.environ["NODUS_INPUT_DATA"], "data.pt"))
model = torch.nn.Linear(16, 1).cuda()
opt = torch.optim.SGD(model.parameters(), lr=0.1)
x, y = d["x"].cuda(), d["y"].cuda()
for step in range(200):
opt.zero_grad()
loss = torch.nn.functional.binary_cross_entropy_with_logits(model(x), y)
loss.backward()
opt.step()
print(f"final loss {loss.item():.4f}")
os.makedirs("/nodus/outputs/model", exist_ok=True)
torch.save(model.state_dict(), "/nodus/outputs/model/model.pt")
outputs:
- name: model
path: /nodus/outputs/model
- name: eval
dependsOn: [train]
inputs:
- name: model
fromStage: train
output: model
template:
resources:
cpu: "2"
memory: 4Gi
command:
- python
- -c
- |
import os, torch
model = torch.nn.Linear(16, 1)
model.load_state_dict(torch.load(os.path.join(os.environ["NODUS_INPUT_MODEL"], "model.pt")))
x = torch.randn(1024, 16)
acc = ((model(x) > 0).float() == (x.sum(dim=1, keepdim=True) > 0).float()).float().mean()
print(f"accuracy {acc.item():.3f}")
Terminal window
$ nodus apply -f pipeline.yaml
pipeline.nodus.dev/train-eval created

The console and the Python SDK also offer this prepare, train and evaluate shape as the train-eval template.

Terminal window
$ nodus get pipeline/train-eval -w
NAME PHASE STAGES COST AGE
train-eval Running 1/3 $0.02 3m
$ nodus get jobs -l nodus.dev/pipeline=train-eval
$ nodus logs job/train-eval-train -f

Each stage runs as a Job named <pipeline>-<stage>, so every Job command works on a stage: logs, exec, describe and cp. status.stages shows each stage’s phase, Job, start and finish times and cost. The pipeline name plus the stage name can be at most 62 characters.

  • dependsOn lists the stages that must succeed first. The stages form a graph with no cycles; a cycle is rejected with the path that closes it.
  • inputs hands an output of an earlier stage to this one: fromStage names the stage (it must also be in dependsOn) and output names one of its declared outputs. The output is mounted read-only at /nodus/inputs/<name>, and NODUS_INPUT_<NAME> holds that path.
  • maxParallel limits how many stages run at once. It defaults to the number of stages, so every ready stage starts.
  • defaults holds Job fields every stage inherits. A field a stage sets wins; objects merge field by field, and a list a stage sets replaces the default list.
failurePolicy What happens
FailFast (default) The running stages are cancelled, the stages not yet started are skipped, and the Pipeline fails
RunIndependent Stages that do not depend on the failed stage run to the end; the stages that depend on it are skipped; the Pipeline then fails

The Pipeline’s status.reason is StageFailed and its message names the failed stages. Each stage keeps its own status.reason, so nodus describe job/<pipeline>-<stage> shows why it failed.

maxCostUSD caps the whole Pipeline: the spend of every stage counts against it. When the Pipeline reaches the cap, every running stage saves its state and suspends, no new stage starts, and the Pipeline becomes Suspended with reason MaxCostReached. Raise the cap to resume:

Terminal window
$ nodus patch pipeline/train-eval --type merge --patch '{"spec":{"maxCostUSD":"1.00"}}'

The cap can only be raised. A stage can set a lower maxCostUSD of its own.

nodus suspend pipeline/x, resume and cancel apply to every stage that has not finished. They set the Pipeline’s spec.state, which its stage Jobs follow. Deleting a Pipeline deletes its stage Jobs and their outputs; ttlSecondsAfterFinished deletes a finished Pipeline for you.