AI Factory

Ship the factory,
not just the model.

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.

In the console

Artifacts

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.

Artifacts
6 total
NameSource runCreated
synthetic-tasksflyte.io/kinddata
u116e4389f847405a ›
6 days ago
promoted-modelflyte.io/kindmodel
u6f7kf6crcq8wpg9kg ›
6 days ago
eval-reportflyte.io/kinddata
u6f7kf6crcq8wpg9kg ›
6 days ago
inference-endpointflyte.io/kinddata
uwbwvdrsf2gzj27gmv ›
6 days ago
policy-checkpointflyte.io/kindmodel
uwbwvdrsf2gzj27gmv ›
6 days ago
rl-tasks-datasetflyte.io/kinddata
uwbwvdrsf2gzj27gmv ›
6 days ago
No artifacts match — the factory hasn't made one of those yet.
In the console

Multi-node distributed training

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.

train_rl_clustered ● Succeeded Duration 5h 39m Actions 1.3k
SummaryLogsMetricsReportsCode
GRPO + LoRA training progress ✓ complete
Base model Qwen/Qwen3-8B · Iterations 20/20 · Group size 6 · Prompts/iter 32 · LoRA rank 16 · LR 2e-05
1.027
▼ −0.023
Mean reward
83.3%
▲ +0.3%
Accuracy
85.9%
▲ +25.0% vs base
Eval accuracy (held-out)
96.9%
▼ −0.1%
Format rate
v20
latest
Adapter version
Reward & correctness vs. iteration mean reward accuracy format rate
iteration 20/20
In the console

Artifact lineage

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.

In the console

Evals per checkpoint

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.

eval_and_promote ● Succeeded Duration 4m 3s Trigger eval-on-new-checkpoint
Evaluation: candidate vs base candidate
32
held-out
Eval tasks
43.75%
pass@1
Candidate
37.50%
pass@1
Base (promoted)
+6.25%
candidate − base
Delta
PASS
delta ≥ 0
Auto gate
taskdifficultycandidatebase
Filter_36434_leasypasspass
Leetcode_32641_lhardfailfail
Filter_4304_leasypasspass
Leetcode_590_lmediumpassfail
Algorithm_20713_leasypasspass
Leetcode_44191_lhardfailfail
Taco_88948_leasypasspass
Data_Structure_11104_leasypasspass
In the console

OnArtifact triggers

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.

Triggers
5 of 5 active
StatusTriggerRuns taskLast run
merge-synthetic-tasks
On new version of synthetic-tasks
de-cpu.merge_synthetic_into_dataset
6 days ago
nightly-synthetic-generation
At 06:00 AM (UTC) — next run in 8h
de-cpu.nightly_synthetic_batch
2 days ago
serve-new-checkpoint
On new version of policy-checkpoint
inference-ops.refresh_inference_service
6 days ago
train-on-new-dataset
On new version of rl-tasks-dataset
trainer.train_grpo
6 days ago
eval-on-new-checkpoint
On new version of policy-checkpoint
eval-cpu.eval_and_promote
6 days ago
The model you own, serving

Serving Traffic
You Own.

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.

score_batch.py
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")
serve_llm.py
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")
review_app.py
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)),)
Wired with events

The Whole Factory
Is a Trigger Away.

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.

dataset.py
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
01

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.

02

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.

03

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.

trainer.py
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
01

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.

02

trigger=flyte.OnArtifact(...)

A new dataset version starts training by itself. Events, not polling — and no scheduler between your teams.

03

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.

04

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.

evals.py
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
01

reusable=flyte.ReusePolicy(...)

Containers stay warm between checkpoints — the model loads once and 10 concurrent tasks share each replica. Zero cold starts per eval.

02

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.

03

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.

serve.py
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
01

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.

02

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.

03

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.

04

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

Enterprise grade Flyte

Open source at the core.

Union is built on Flyte, the open-source AI runtime we create and maintain under the Linux Foundation AI & Data.

A Linux Foundation AI & Data Project
4000+ companies using Flyte today
18M+ Flyte SDK downloads

Stand Up Your First Component.

Bring one stage of your pipeline — data, training, or evals. Run it durable, event-wired, and inside your own cloud this week.