Kubeflow is the set of Kubernetes controllers and Python SDKs that teams reach for when GPU work outgrows one notebook and one SSH session. Calling it an ML platform hides the useful part: each piece is a controller that turns a declarative object into pods, and each fails in its own way. This article covers what each controller owns, how a pipeline step becomes a pod, how a distributed training job becomes GPU pods that find each other, and how a queue keeps a shared cluster from deadlocking.
You should know what a Kubernetes Job is and what torchrun does; no prior Kubeflow install is assumed.
What Kubeflow is now
Kubeflow today is a family of subprojects rather than one product. The project's own introduction says the subprojects are designed to be usable on their own as well as together in a distribution. That matters for planning. You can install Kubeflow Trainer on a cluster that has no Pipelines UI, or run Pipelines without Katib. Most production teams install only the parts they use.
The parts that matter for GPU work:
- Kubeflow Pipelines (KFP) compiles a Python function graph into an intermediate representation and runs each step as a container. It records inputs, outputs and artifacts so that runs can be compared and cached.
- Kubeflow Trainer (version 2, which replaces the older Training Operator and its per-framework kinds such as
PyTorchJob) runs distributed training as aTrainJobthat references a reusable runtime. - Katib runs hyperparameter search. An
Experimentcreates many trials, collects a metric from each, and proposes the next set of values. - Notebooks give users browser IDEs that run as pods in their own namespaces.
- The model registry stores model versions and the metadata that links a version back to the run that produced it.
Serving with KServe and features with Feast sit next to this family as ecosystem projects. Queueing is not part of Kubeflow at all. Kueue is a Kubernetes SIG project, and Trainer integrates with it. That split is the first lesson. Kubeflow decides what pods a job needs, while Kueue and the Kubernetes scheduler decide when and where they run.
The architecture: controllers all the way down
Figure 1 shows the shape. A user talks to two SDKs. The kfp SDK compiles and submits pipelines to the KFP API server. The kubeflow SDK submits training jobs directly as TrainJob objects. Each controller then creates lower-level objects. KFP's backend hands the run to a workflow engine (Argo Workflows in the standard install), which starts one pod per step. Trainer turns a TrainJob into a JobSet, which is a group of Jobs with a headless Service so the ranks can find rank 0. Katib creates trial objects, which are often TrainJobs themselves.
Durable state lives outside the controllers: the object store for artifacts and checkpoints, the metadata database behind KFP and the model registry, and the container registry. When a GPU pod dies, only what it wrote to the object store survives. Design every step on that assumption.
Multi-user installs map each team to a namespace (a profile), so GPU quota becomes a namespace concern, which is what Kueue handles.
Pipelines: from Python to pods
A KFP pipeline is ordinary Python with two decorators. @dsl.component marks a function that will run in its own container. @dsl.pipeline wires components together by passing one task's outputs to another's inputs. The compiler walks that graph and writes an intermediate-representation YAML file holding a PipelineSpec. The YAML, not your Python, is what the server runs, so you can version it and submit it from CI.
from kfp import compiler, dsl
from kfp.dsl import Dataset, Input, Model, Output
@dsl.component(base_image="python:3.11", packages_to_install=["pandas", "pyarrow"])
def prepare(raw_uri: str, train: Output[Dataset]):
import pandas as pd
df = pd.read_parquet(raw_uri)
df = df.dropna(subset=["text"])
df.to_parquet(train.path) # KFP uploads train.path to the artifact store
@dsl.component(base_image="registry.example.com/ft-trainer:2026.10")
def finetune(train: Input[Dataset], epochs: int, model: Output[Model]):
import subprocess
subprocess.run(["python", "/app/train.py", "--data", train.path,
"--epochs", str(epochs), "--out", model.path], check=True)
@dsl.pipeline(name="nightly-finetune")
def nightly(raw_uri: str, epochs: int = 2):
prep = prepare(raw_uri=raw_uri)
ft = finetune(train=prep.outputs["train"], epochs=epochs)
ft.set_accelerator_type("nvidia.com/gpu").set_accelerator_limit(1)
ft.set_retry(num_retries=2)
ft.set_caching_options(False) # never reuse a stale model
compiler.Compiler().compile(nightly, package_path="nightly.yaml")Three details in that code carry most of the operational weight.
- Artifacts are paths, not objects.
Output[Dataset]gives the step a local path. The launcher copies the file to the artifact store after the step exits, and the next step'sInput[Dataset]is downloaded before it starts. A 200 GB dataset passed this way is copied twice. Pass large data as a URI string parameter and let the training code stream it. - Caching is keyed on the component and its inputs. If neither changed, KFP skips the step and reuses the old outputs. That is a gift for data preparation and a trap for anything whose real input is "whatever is in the bucket today". Disable it on such steps, as above, or pass a date parameter so the key changes.
- A GPU step is a single pod.
set_accelerator_limitasks for GPUs on one node. A step cannot span nodes. For multi-node training, the step should submit aTrainJoband wait for it, which is the subject of the next section.
Trainer v2: TrainJobs and runtimes
Trainer v2 splits a training job into two objects with two owners. A platform team writes ClusterTrainingRuntime (cluster-wide) or TrainingRuntime (one namespace) objects. These are templates that describe the JobSet shape, the launcher (for example torchrun for the torch-distributed runtime), default images and any initializer steps. A practitioner writes a small TrainJob that names a runtime and overrides only what differs: image, command, node count and per-node resources.
apiVersion: trainer.kubeflow.org/v1alpha1
kind: TrainJob
metadata:
name: llama-ft-0412
namespace: team-nlp
labels:
kueue.x-k8s.io/queue-name: team-nlp-queue
spec:
runtimeRef:
name: torch-distributed
kind: ClusterTrainingRuntime
trainer:
image: registry.example.com/ft-trainer:2026.10
command: ["python", "/app/train.py", "--ckpt-dir", "s3://ckpt/llama-ft-0412"]
numNodes: 2
resourcesPerNode:
requests:
cpu: "32"
memory: "200Gi"
nvidia.com/gpu: "8"The controller merges this with the runtime and creates a JobSet with two pods of eight GPUs each. The runtime's launcher starts one process per GPU and supplies the rendezvous settings that torch.distributed needs, so the training script calls init_process_group exactly as it would under torchrun on bare metal. The same job can be submitted from Python, which is convenient from a notebook or a KFP step:
from kubeflow.trainer import CustomTrainer, TrainerClient
def train_fn():
import torch.distributed as dist
dist.init_process_group("nccl")
... # build model, wrap with FSDP, train, save checkpoints
client = TrainerClient()
job_id = client.train(
trainer=CustomTrainer(
func=train_fn,
num_nodes=2,
resources_per_node={"cpu": 32, "memory": "200Gi", "gpu": 8},
)
)
for line in client.get_job_logs(job_id, follow=True):
print(line)Splitting the object in two is the main design win of v2. Team-wide settings (NCCL environment, shared-memory volume, GPU tolerations, launcher) live in one reviewed runtime; per-experiment settings live in a dozen lines of YAML. Trainer v2 is still a v1alpha1 API, so pin the controller version and read release notes before upgrading.
Queueing GPUs with Kueue
A distributed training job is useless until all of its pods run. If two 16-GPU jobs meet 24 free GPUs and pods start one at a time, each can grab 12 and wait forever for the rest, holding GPUs idle. Plain Kubernetes does nothing to prevent this partial-admission deadlock.
Kueue prevents it by gating whole workloads. A TrainJob carrying the kueue.x-k8s.io/queue-name label is created suspended. Kueue adds up the resources of all its pods and checks them against the ClusterQueue quota behind that LocalQueue. It unsuspends the job only when the whole job fits. Quotas are written per resource flavor, for example one flavor per GPU type, and teams in a cohort can borrow each other's idle quota. Preemption policies decide whether a borrowed slot can be taken back.
A practical layout is one ClusterQueue per team in a shared cohort, plus a low-priority queue for Katib trials that may be preempted. Training code must then survive preemption by checkpointing often enough that losing the newest checkpoint is cheap.
Katib: the tuning loop
Katib automates the outer loop of tuning. An Experiment declares an objective metric, a search space, an algorithm and a trial template. The controller asks a suggestion service for parameter sets, creates one trial per set from the template, collects the metric each trial reports, and repeats until it hits maxTrialCount or the goal.
apiVersion: kubeflow.org/v1beta1
kind: Experiment
metadata: {name: lr-sweep, namespace: team-nlp}
spec:
objective: {type: minimize, objectiveMetricName: eval_loss}
algorithm: {algorithmName: bayesianoptimization}
parallelTrialCount: 4
maxTrialCount: 24
maxFailedTrialCount: 4
parameters:
- name: lr
parameterType: double
feasibleSpace: {min: "1e-5", max: "3e-4"}
# trialTemplate: a TrainJob or Job spec using ${trialParameters.lr}Two settings decide what this costs. parallelTrialCount times the GPUs per trial is the peak GPU demand, so size it against the queue's quota, not the cluster. maxFailedTrialCount stops a broken template from burning a night of GPU time on trials that crash at import. A misnamed metric looks like trials with no result, so start with maxTrialCount of 2 to check the plumbing.
Worked example: a nightly fine-tune
Take a team that fine-tunes a 7B-parameter model every night on fresh support tickets. The cluster has four 8-GPU nodes, and the team's queue has a nominal quota of 16 GPUs. The numbers below are illustrative. Measure your own.
- At 01:00 a recurring KFP run starts. The
preparestep (CPU only) reads yesterday's tickets, filters and tokenizes them, and writes about 3 GB of shards to the object store. It returns their URI as a string, not as a dataset artifact, so nothing is copied twice. - Unlike the single-GPU sketch above, this pipeline's
launchstep submits a two-nodeTrainJobwith 8 GPUs per node and the URI as an argument. Kueue sees 16 GPUs requested against 16 of quota. If another job holds 8 of them, the TrainJob waits suspended rather than half-starting. - Once admitted, 16 ranks start. Each streams its shard range, trains with FSDP, and rank 0 writes a sharded checkpoint every 20 minutes. A run takes about 2 hours.
- The
launchstep polls the job until it completes, then anevaluatestep on one GPU scores the final checkpoint against a fixed test set. - If the score beats the current production version, a
registerstep records the new model version with the run ID, dataset URI and score. Otherwise the run ends and the artifacts expire under the bucket's lifecycle rule.
Every step is idempotent given its inputs, so a retry repeats work but never corrupts state. That property, more than any Kubeflow feature, makes the run safe to leave alone.
Failure modes
These are the failures teams actually hit, with the signal each one gives and the fix.
| Failure | What you see | Fix |
|---|---|---|
| Partial admission without a queue | Jobs hold GPUs and sit at 0% utilization; pods stay Pending | Route every TrainJob through Kueue; never let multi-pod GPU jobs bypass it |
| Rendezvous timeout | Ranks log connection errors to rank 0 and the job fails after the timeout | Check the headless Service DNS and NetworkPolicy between pods; confirm all pods were admitted together |
| NCCL hangs or slow all-reduce | Steps stall with GPUs at 100% and no progress, or throughput far below a single node | Run an NCCL test in the same runtime; check the fabric interface settings and the shared-memory volume |
| Stale cached step | A run finishes in seconds with yesterday's model | Disable caching on steps that read changing external data, or pass a date parameter |
| Preempted without checkpoints | Hours of training lost when a higher-priority job arrives | Checkpoint on a time interval and resume from the newest checkpoint at start-up |
Operating it and the trade-offs
Install only what you run. Every controller is another upgrade path and another webhook that can block pod creation when unhealthy.
Pin images by digest and runtimes by version. A pipeline that pulls :latest cannot be reproduced.
Watch the queue, not just the GPUs. The useful dashboards are pending workloads per queue, time to admission and GPU allocation versus actual utilization. A GPU that is allocated and idle is the most expensive failure in the cluster, and it never raises an alert by itself.
Know the trade-offs. Kubeflow gives you a reproducible graph, lineage and a shared queue at the price of several controllers, an object store, a database and an auth layer to operate. Slurm is simpler for a pure training cluster, Ray suits dynamic Python better than a static graph, and managed pipeline services trade operations work for lock-in. Kubeflow fits best where Kubernetes is already the platform and several teams share GPUs.
For related depth, see PyTorch FSDP for the training loop inside each rank, NCCL collectives for what the ranks do on the wire, hyperparameter search strategies for what Katib's algorithms are doing, and MLflow for GPU training ops if you track experiments outside Kubeflow.
What to do next
- Draw your current GPU workflow as steps and mark which steps need more than one node. Only those need Trainer; the rest can be plain pipeline steps.
- Install Kueue first, with one ClusterQueue per team in a shared cohort, and require the queue label on every GPU workload.
- Install Trainer, write one ClusterTrainingRuntime for PyTorch with your NCCL settings, tolerations and shared-memory volume, and run a two-node smoke job.
- Convert one existing training script to a TrainJob and confirm it resumes from its newest checkpoint after you delete a pod.
- Build a KFP pipeline around it, passing large data as URIs and disabling caching on steps that read changing data.
- Add a Katib experiment with two trials to prove metric collection, then raise the counts.
- Add dashboards for queue wait time and allocated-but-idle GPUs before inviting more teams.