Skip to content

Runs: pipelines that finish

A flow that automates a house never ends — a value arrives, nodes fire, and it waits for the next one. A research pipeline is the other shape: parameters go in, stages execute in order, and at some point it is done and has produced something worth keeping. Fluksio calls the second one a run, and it is the same engine either way.

This is what makes Fluksio usable where Kedro, MLflow or ClearML would be: a run has parameters that identify it, a result, per-step metrics, artifacts and a place in a queryable history — without a second server, and without paying a project bootstrap on every execution.

A batch flow

Set mode: "batch" on the flow and name the messages its result should hold:

{
  "name": "train_polymer_gnn",
  "mode": "batch",
  "outputs": ["final_loss", "report"],
  "inputs": [
    {"spec": {"name": "lr", "dtype": "float"}, "initial": 0.1},
    {"spec": {"name": "steps", "dtype": "int"}, "initial": 3000}
  ],
  "nodes": [{"id": "train", "timeout": 7200, "device": "gpu", "...": "..."}]
}

A batch flow is built and validated like any other — it appears on the canvas, its ports are type-checked — but it is never activated: no subscriptions, no schedules, no webhooks. It runs when a run asks it to, and not otherwise.

Its inputs are its parameters. A run supplies values for them; anything it does not supply keeps the declared initial value.

One thing a batch flow may not do is rate-limit a port (interval). A rate limit holds a value back for a timer to release, and a run has no timer — the value would be dropped rather than delayed, so submitting is refused instead.

Submitting

curl -X POST $FLUKSIO/runs/flows/train_polymer_gnn \
     -H "Authorization: Bearer $TOKEN" -H 'Content-Type: application/json' \
     -d '{"params": {"lr": 0.3, "steps": 4000}, "seed": 7}'

The answer is immediate and the run is queued; a training run is measured in hours, so nothing waits for it. Poll GET /api/v1/runs/{id} for its status, result, per-node record and artifacts.

Wrong parameters are refused before anything executes — an undeclared name, or a value of the wrong type, comes back as a 422 naming the problem.

Sweeps

An ensemble is the same parameters at different seeds; a grid search is the parameters spread out. Both are one call, and the caller builds the list:

curl -X POST $FLUKSIO/runs/flows/train_polymer_gnn/sweep \
     -H "Authorization: Bearer $TOKEN" -H 'Content-Type: application/json' -d '{
  "runs": [{"params": {"lr": 0.1}, "seed": 1}, {"params": {"lr": 0.3}, "seed": 1}]
}'

They share a group_id, so GET /api/v1/runs?group=… is the sweep, and they execute in parallel. That is safe because each run has a state backend of its own: message names are global keys, so two runs of one flow would otherwise overwrite each other's values. They do not.

Producing values before you are finished

A training loop has numbers worth keeping long before it has a result. Those numbers are outputs, not logs: a node declares a port for them and produces them over time, which in Python is a generator.

def process(lr, steps):
    loss = 1.0
    for _ in range(steps):
        loss = train_one_step(lr)
        yield {"loss": loss}          # published now, on the `loss` port
    return {
        "weights": fluksio.save_artifact(dump(model), "weights.npz"),
        "final_loss": loss,
    }

Mark the port it streams on, so the flow says what it does:

{"name": "loss", "dtype": "float", "stream": true}

Every yield is published the instant it happens — same port, same type check, same place on the canvas as any other value. Whatever the generator returns is the node's result, and is what downstream nodes read. If you never return, the last thing you yield is the result instead.

This is the whole reason the framework does not have a logging API. A metric that escapes through log_metric() is undeclared: invisible to validation, absent from the canvas, and stored somewhere the graph knows nothing about. A metric that leaves through a port is a message — so a chart binds to it directly, a downstream node can consume it, and the run keeps its series without anyone asking.

Where a yield cannot reach — the value comes from inside somebody else's callback, and they call you rather than the other way round — fluksio.emit writes the same ports the same way:

import fluksio


def process():
    model.fit(callbacks=[LambdaCallback(
        on_epoch_end=lambda epoch, logs: fluksio.emit(loss=logs["loss"])
    )])
    return {"weights": ...}

What a run does with them

Every number a node emits is kept as the run's series, stepped by the count of emissions on that message. Read one back with GET /api/v1/runs/{id}/metrics?name=<flow>.loss, or compare runs:

GET /api/v1/runs/series/compare?ids=<a>,<b>,<c>&metric=<flow>.loss

That answers in the series shape a chart widget already draws, so three training curves side by side is a widget binding. During a run the values also arrive live on the flow socket, so a chart bound to the port fills in as the training goes.

A streaming port may set interval to thin out what reaches the canvas — the run's history still keeps every value, because the interval is asking for the display not to be flooded, not for the curve to have holes in it.

Emitting has a second effect: a node's timeout measures silence, not duration. A node that yields every few seconds can run for hours under a timeout of 300; one that says nothing for longer than its timeout is killed. Set timeout on a long node to how long it may plausibly go quiet.

In a live flow, an emission also wakes whatever is downstream of it, exactly as a subscriber publishing does. In a run it does not: a run's graph is scheduled once, and three thousand mid-node cascades would leave "the run has finished" with nothing to mean.

Artifacts

Bytes never travel as a message. save_artifact writes them to a content-addressed store and returns a small reference — digest, size, media type, name — which is what an artifact-typed port carries:

def process(weights):          # requires: weights, dtype "artifact"
    path = fluksio.load_artifact(weights)
    ...

Because the address is the content's hash, a sweep whose fifty configs share one preprocessed input stores it once, and a reference stays valid wherever the store is reachable from. Artifacts a run produced are listed on it and downloadable at GET /api/v1/artifacts/{digest}.

Objects that cannot be serialized

A live model, a DataLoader, a JAX-compiled function — these do not cross a node boundary, and no framework flag will make them. There are exactly two patterns, and they are both deliberate:

  • Keep them in one node. Stages that must share live memory are one node. Building the model and training it is one stage; the fact that Kedro would make them two nodes is Kedro's problem, not a structure worth reproducing.
  • Cross at a checkpoint. Save what matters as an artifact and rebuild from it on the other side. That is the boundary that also survives the next node running on a different machine.

Running a node somewhere else

A node that needs a GPU declares the label of the machine that has one:

{"id": "train", "device": "gpu", "device_policy": "require", "timeout": 7200}

A worker on that machine dials out to the engine, because the engine generally cannot reach it — different network, no inbound route — and because nothing should expose Redis across hosts. Install it on the box, mint it a token, and start it:

curl -X POST $FLUKSIO/workers/tokens -d '{"name": "gpu-dev"}'   # once, as an admin

pip install fluksio-worker

fluksio-worker \
  --url wss://api.example.com/api/v1/workers/attach \
  --token "$FLUKSIO_WORKER_TOKEN" \
  --labels gpu,cuda12 \
  --python /opt/torch-venv/bin/python

fluksio-worker is its own distribution — the agent, the node runner, and websockets. Nothing of the engine, so a GPU box does not install a database driver to run a training step. Where pip is not an option, the two files still work copied into one directory and run with python agent.py …; the engine serves the runner at GET /api/v1/workers/runtime.

--python is the interpreter node code runs on, which is how the GPU box keeps its CUDA wheels without the engine ever installing them. The node's source travels with every call, so nothing has to be deployed there.

A few consequences worth knowing:

  • import fluksio inside a node is the worker's own reporter — emit, save_artifact, load_artifact — installed before the node's code runs, so the installed fluksio package (if the box has one) never shadows it.
  • A node bound to a device is compiled on that machine. A node importing torch is correct on the GPU box and a missing module on the engine, so checking it on the engine would fail a node that is fine.
  • If nothing carrying the label is attached, the run stays queued and says what it is waiting for. Submit first, switch the GPU box on later.
  • Cancelling a run kills what it is executing, there or here, and leaves other runs of the same node alone.
  • If the worker disappears mid-call, the run fails in seconds with worker went away mid-call rather than waiting out its timeout.
  • device_policy: "prefer" runs locally when no such worker is attached; "require" (the default) waits for one.

Durability

Submitting journals the run to a Redis stream of its own, separate from the one the automations use — a burst of five hundred sweep runs must not stand between a house and its heating. An engine that is down when a run is submitted picks it up when it starts.

From the moment a run is claimed, its database row is the record and the queue is finished with it. Redelivering two hours of training because an acknowledgement was late is not recovery; instead a running run refreshes a lease, and one whose lease goes stale is marked abandoned — which is what a run whose engine was killed mid-training becomes.

What this costs, compared

The repository ships a benchmark that measures submitting a run against a Kedro project doing the same nothing:

fluksio — submit accepted     median   15.3 ms
          submit -> result    median   60.8 ms
kedro   — kedro run           median 1109.5 ms

The difference is not the orchestration; it is that Fluksio does not boot a project per run. The engine is already up, and the workers already have the node's code compiled. On a 510-run sweep, that gap is about nine minutes of pure startup that never happens.

See also