A frontier lab's edge is the loop: data refined, checkpoints trained, evals scored, the next run already queued. Union runs that loop end to end, durable by default and wired together with artifacts and events, entirely inside your perimeter.
Six parts of the factory. Follow a card to the detail.
Every dataset, checkpoint, and report in the factory is a named, versioned artifact — typed, searchable, and linked to the run that produced it. Try the search.
A GRPO + LoRA run fanned out across nodes — 1.3k actions in one run, with live metrics, custom reports, and reruns from any step. Replay the run or scrub through its iterations.
The assembly line, inspectable: which run produced each artifact version and which components consume it downstream — the wiring diagram is live data, not a slide. Click a node to trace it.
Every checkpoint version fires an eval run automatically. The promotion gate reads the report — promoted-model only moves when the numbers clear the bar, on the record. Switch candidates to watch the gate decide.
| task | difficulty | candidate | base |
|---|---|---|---|
| Filter_36434_l | easy | pass | pass |
| Leetcode_32641_l | hard | fail | fail |
| Filter_4304_l | easy | pass | pass |
| Leetcode_590_l | medium | pass | fail |
| Algorithm_20713_l | easy | pass | pass |
| Leetcode_44191_l | hard | fail | fail |
| Taco_88948_l | easy | pass | pass |
| Data_Structure_11104_l | easy | pass | pass |
The event wiring from the code above, as operators see it: on a new version of policy-checkpoint, run eval-and-promote. Flip a toggle to take a component offline — no code changes.
A factory that can’t ship isn’t a factory. Batch and real-time inference run on the same runtime that trained the model, against the same artifacts, inside the same perimeter — so the checkpoint that passed the eval gate is the checkpoint answering requests, with the lineage to prove it.
import flyte infer = flyte.TaskEnvironment( name="batch-infer", resources=flyte.Resources(cpu="2", memory="8Gi", gpu="L40s:1"), interruptible=True, # spot by default) @infer.task(trigger=flyte.OnArtifact("promoted-model"), cache="auto")async def score_all(rows: flyte.io.Dir) -> flyte.io.File: # a promoted checkpoint lands, scoring starts itself shards = await flyte.map(score, rows, concurrency=2_000) # every element checkpointed: a reclaimed node costs one element return publish(merge(shards), name="scored-rows")
from flyteplugins.vllm import VLLMAppEnvironmentimport flyte vllm_app = VLLMAppEnvironment( name="qwen3-serving", model_hf_path="Qwen/Qwen3-8B", model_id="qwen3-8b", resources=flyte.Resources( cpu="4", memory="16Gi", gpu="L40s:1", disk="10Gi"), scaling=flyte.app.Scaling( replicas=(0, 2), scaledown_after=300), stream_model=True, # weights stream from blob store to GPU) app = flyte.serve(vllm_app)print(f"OpenAI-compatible API: {app.url}/v1")
from flyteplugins.fastapi import FastAPIAppEnvironmentfrom fastapi import FastAPIimport flyte api = FastAPI() @api.get("/eval/{version}")async def eval_report(version: str): return load_report(version) # reads the eval-report artifact env = FastAPIAppEnvironment( name="eval-review", app=api, resources=flyte.Resources(cpu="1", memory="1Gi"), scaling=flyte.app.Scaling(replicas=(1, 4)),)
The teams already running the factory, and the number each of them reported.
No scheduler glue, no polling loops, no orchestration YAML. A component is a Python task; the wiring between components is an artifact event. Walk the line: data, training, evals.
import flyte de = flyte.TaskEnvironment(name="data-engine", resources=flyte.Resources(cpu=8)) @de.task( trigger=flyte.OnArtifact("production-traces"), cache="auto",)async def build_dataset(traces: flyte.io.Dir) -> flyte.io.File: # fresh production traces land → the data component wakes itself tasks = await flyte.map(synthesize, traces, concurrency=2_000) # 2,000 spot containers; each element checkpointed, nothing rebuilt ds = merge_and_filter(tasks) return publish(ds, name="rl-tasks-dataset") # └── v15 created → the foundry fires
trigger=flyte.OnArtifact(...)
Production traces are an artifact too. A new version lands and the dataset build starts itself — no scheduler between serving and data.
flyte.map(..., concurrency=2_000)
Durable fan-out across thousands of spot containers. Every element is checkpointed — a reclaimed node costs one element, not the run.
publish(..., "rl-tasks-dataset")
The dataset becomes a named, versioned artifact. The training component subscribes to it — the teams never import each other’s code.
import flyteimport flyte.errorsfrom flyte.clustered import ClusteredTaskEnvironment, TorchRun trainer = ClusteredTaskEnvironment( name="trainer", resources=flyte.Resources(gpu="A100:4"), replicas=4, nproc_per_node=4, # 4 nodes × 4 workers runtime=TorchRun(rdzv_backend="static"),) @trainer.task( trigger=flyte.OnArtifact("rl-tasks-dataset"), retries=3, cache="auto",)async def train_grpo(dataset: flyte.io.File) -> flyte.io.Dir: # a new dataset version lands → 16 workers start themselves try: ckpt = await grpo_epoch(dataset) except flyte.errors.OOMError: # infra failures are control flow: same step, bigger boxes ckpt = await grpo_epoch.override( resources=flyte.Resources(gpu="H100:8") )(dataset) return publish(ckpt, name="policy-checkpoint") # └── v10 created → eval + serving components fire
ClusteredTaskEnvironment(replicas=4)
One environment declaration turns the task into a torchrun cluster — 4 nodes, 16 workers, rendezvous wired for you. Multi-node is config, not infrastructure.
trigger=flyte.OnArtifact(...)
A new dataset version starts training by itself. Events, not polling — and no scheduler between your teams.
except flyte.errors.OOMError
Infrastructure failures are typed exceptions. The handler re-runs the same step on bigger hardware, mid-run — a try block that reaches the cluster.
publish(..., "policy-checkpoint")
The checkpoint outlives the run as a versioned artifact. Eval and serving subscribe to it — they fire on v10 without ever importing this module.
import asyncioimport flytefrom flyte.extras import TokenBatcher evals = flyte.TaskEnvironment( name="evals", resources=flyte.Resources(gpu="A10G:1"), reusable=flyte.ReusePolicy(replicas=2, concurrency=10),) @alru_cache(maxsize=1)async def get_batcher() -> TokenBatcher: # one batcher per container — every task feeds the same GPU queue batcher = TokenBatcher(await load_vllm(), target_batch_tokens=32_000) await batcher.start() return batcher @evals.task(trigger=flyte.OnArtifact("policy-checkpoint"))async def eval_checkpoint(ckpt: flyte.io.Dir) -> flyte.io.File: batcher = await get_batcher() futures = [await batcher.submit(t) for t in held_out_tasks()] report = score(await asyncio.gather(*futures)) return publish(report, name="eval-report") # └── every checkpoint scored — the promotion gate reads this
reusable=flyte.ReusePolicy(...)
Containers stay warm between checkpoints — the model loads once and 10 concurrent tasks share each replica. Zero cold starts per eval.
TokenBatcher(target_batch_tokens=...)
Every concurrent task feeds one shared dynamic batcher, so the GPU always has a full queue. Throughput scales without touching the eval code.
publish(..., "eval-report")
Every checkpoint version fires this eval and the report goes on the record. Promotion only moves when the numbers clear the bar.
from flyteplugins.vllm import VLLMAppEnvironmentimport flyte, flyte.app app = VLLMAppEnvironment( name="policy-endpoint", model_id="policy-v6", model_path=flyte.app.RunOutput( type="directory", run_name=train_run.name, ), resources=flyte.Resources(gpu="H100:2", memory="64Gi"), stream_model=True, scaling=flyte.app.Scaling( replicas=(0, 8), scaledown_after=300, ),) endpoint = flyte.serve(app)# OpenAI-compatible API at {endpoint.url}/v1 — inside your VPC
VLLMAppEnvironment(...)
The last component is a deployment, declared the same way every other component is. An OpenAI-compatible endpoint, behind your own load balancer, reachable only from inside your network.
model_path=RunOutput(...)
The endpoint is fed by the training run's output directly — no manual weight copy, no separate model registry to keep in sync. The artifact that passed the eval gate is the artifact that serves.
stream_model=True
Weights stream from your object store straight to GPU memory rather than landing on local disk first, which is what keeps a cold start from being a coffee break.
Scaling(replicas=(0, 8))
Scale to zero when nobody is asking, up to eight replicas when they are. The GPUs go back to the training queue in between — same fleet, same scheduler.
From the reference implementation: github.com/unionai-oss/model-factory
Not a diagram — running code. Two factories, CI-tested and deployable to any Union cluster.
GRPO fine-tuning with sandboxed unit tests as verifiable rewards. Four teams, six artifacts, zero shared code — every hand-off is an artifact event.
A small language model fine-tuned to right-size flyte.Resources requests — trained in a simulator, verified with real cluster episodes. The factory improving the factory.
Union is built on Flyte, the open-source AI runtime we create and maintain under the Linux Foundation AI & Data.
Every component above assumes three things underneath it: a perimeter the work never leaves, a runtime that survives the hardware, and a scheduler that decides who gets the GPUs.
Every component above runs inside your own perimeter. The control plane holds references, never payloads, and your data never transits ours.
See the architecture →A factory only runs if the line survives. Recover, fork, and replay any run, and change the hardware inside an except block.
Explore the runtime →The queues that feed the factory: priority, quotas, and gang scheduling across clusters, clouds, and accelerator classes.
See the scheduler →Bring one stage of your pipeline — data, training, or evals. Run it durable, event-wired, and inside your own cloud this week.