Kubeflow Pipelines and TFJob: When the Extra Controller Earns Its Keep
Wrapping a training step in a TFJob adds a controller, a CRD, and an RBAC surface. Here is how to tell whether your pipeline actually needs it.
29 Aug 2025, 09:05 UTC

Most model-training pipeline steps start life as a plain container: one image, one command, one pod. That works until someone asks for two workers, per-replica restart behaviour, and a job-level status you can query. The reflex is to wrap the step in a TFJob. Sometimes that is right. Often it adds a controller, a CRD, and an RBAC surface you did not need.
The decision worth making deliberately: does your training step need the Training Operator to reconcile it, or does it just need to run?
What the two options actually are
A Kubeflow Pipelines (KFP) step is a container. The KFP compiler turns your Python pipeline into a pipeline package, and the KFP backend submits it to Argo Workflows, which runs each step as a pod. Parameters and artifacts move between steps through the KFP runtime. Nothing extra is installed beyond KFP itself.
A TFJob is a Kubernetes custom resource in the kubeflow.org/v1 API group. It is not run by Argo. A separate controller — the Training Operator — watches TFJob objects and creates the worker and parameter-server pods described in the spec, then keeps them in the state you asked for. That reconciliation loop is the whole point: if a worker pod dies and the replica spec says restartPolicy: OnFailure, the controller recreates it without the pipeline knowing.
So the real question is whether you want a controller owning the training pods, or whether Argo's pod-level retry is enough.
A worked example: submitting a TFJob from a pipeline step
The pattern is two steps. The first creates the TFJob custom resource; the second polls its status and fails the pipeline if the job fails. Run all of this from a machine with the KFP SDK installed and a kubeconfig pointing at the cluster.
Here is the equivalent TFJob manifest, useful for understanding the shape and for applying by hand when debugging:
apiVersion: kubeflow.org/v1
kind: TFJob
metadata:
name: mnist-trainer
namespace: kubeflow-user-example-com
spec:
tfReplicaSpecs:
Worker:
replicas: 2
restartPolicy: OnFailure
template:
spec:
containers:
- name: tensorflow
image: registry.example.com/mnist-trainer:1.0
command: ["python", "/app/train.py"]
resources:
requests: {cpu: "2", memory: "4Gi", nvidia.com/gpu: "1"}
limits: {cpu: "2", memory: "4Gi", nvidia.com/gpu: "1"}
Confirm the field names against your installed CRD before trusting them — Training Operator releases have changed spec shape over time. On the cluster, kubectl explain tfjob.spec prints the schema the API server actually serves.
The pipeline component builds that same object and posts it through the Kubernetes API:
from kfp import dsl
@dsl.component(base_image="python:3.11-slim", packages_to_install=["kubernetes"])
def submit_tfjob(namespace: str, job_name: str, image: str, workers: int) -> str:
from kubernetes import client, config
config.load_incluster_config()
body = {
"apiVersion": "kubeflow.org/v1",
"kind": "TFJob",
"metadata": {"name": job_name, "namespace": namespace},
"spec": {"tfReplicaSpecs": {"Worker": {
"replicas": workers,
"restartPolicy": "OnFailure",
"template": {"spec": {"containers": [{
"name": "tensorflow",
"image": image,
"command": ["python", "/app/train.py"],
"resources": {
"requests": {"cpu": "2", "memory": "4Gi"},
"limits": {"cpu": "2", "memory": "4Gi"},
},
}]}},
}}},
}
client.CustomObjectsApi().create_namespaced_custom_object(
group="kubeflow.org", version="v1",
namespace=namespace, plural="tfjobs", body=body,
)
return job_name
The wait component polls status.conditions, which the Training Operator populates with types such as Running, Succeeded, and Failed:
@dsl.component(base_image="python:3.11-slim", packages_to_install=["kubernetes"])
def wait_for_tfjob(namespace: str, job_name: str) -> None:
import time
from kubernetes import client, config
config.load_incluster_config()
api = client.CustomObjectsApi()
while True:
obj = api.get_namespaced_custom_object(
group="kubeflow.org", version="v1",
namespace=namespace, plural="tfjobs", name=job_name)
for cond in obj.get("status", {}).get("conditions", []):
if cond.get("status") != "True":
continue
if cond["type"] == "Succeeded":
return
if cond["type"] == "Failed":
raise RuntimeError(f"TFJob {job_name} failed: {cond.get('reason')}")
time.sleep(15)
Two operational details decide whether this works in practice. First, the pipeline pod's service account needs RBAC to create and read tfjobs in the target namespace — without it the submit step fails on a 403, not on anything TensorFlow-related. Second, disable KFP's execution caching on the submit step (the setter name differs between SDK versions; check the SDK reference for yours). Otherwise a re-run with identical inputs can skip creation while the wait step still polls for a job that no longer exists.
What you give up by going through TFJob
| Concern | Plain KFP step | TFJob step |
|---|---|---|
| Artifact passing | Native, typed | Manual, via shared object store or PVC |
| Retry semantics | Argo retries the whole pod | Controller restarts individual replicas |
| Multi-replica training | Not supported | Supported |
| Debugging surface | One pod, one log stream | Controller plus N worker pods |
The artifact row is the one that bites. A TFJob writes to whatever storage your training code is configured for; it does not hand a KFP artifact back to the pipeline. You either mount a shared volume or have the wait step read results from S3-compatible storage and emit them as an artifact explicitly.
One more limitation worth stating plainly: TFJob does not by itself guarantee that all workers start together. Without a gang-scheduling plugin, some replicas can be scheduled while others wait, which is exactly the failure mode that exhausts node resources on large jobs. Always set requests and limits, and treat GPU limits as mandatory rather than nice-to-have.
Checking that it actually worked
After a run, verify at three levels rather than trusting the UI's green checkmark:
kubectl get tfjobs -n kubeflow-user-example-com— the job should exist and show a state;kubectl describe tfjob mnist-trainer -n kubeflow-user-example-comshows replica-level events.- The Kubeflow Pipelines UI Runs view — the wait step's log should show the terminal condition it observed, and the run's lineage graph should link the two steps.
- The artifact store configured for your Kubeflow deployment (commonly MinIO) — confirm the model or metrics object you expect is actually there, since a successful TFJob status does not imply your training code wrote anything.
If the TFJob never appears, the problem is almost always RBAC or a CRD that does not match your Kubeflow release. Check kubectl api-resources | grep tfjobs first; if the resource is absent, the Training Operator is not installed or not the version your pipeline assumes.
The short version: reach for a plain KFP step until you genuinely need multiple replicas with independent restart behaviour. When you do cross that line, budget for the RBAC, the shared artifact path, and the extra debugging surface — those, not the training code, are what usually costs the afternoon.
0 replies
A thoughtful contribution can make all the difference. Be the first to share one.