Union keeps workflow and agent state in durable object storage, not on the node doing the work, so a failed step picks up where it stopped. OOM kills, preempted nodes, and GPU failures stop being incidents.
Try for free
Six properties of the durable runtime, and where each one lives on this page.
Most runtimes hand you a stack trace. Union hands you a typed error and the ability to do something about it. Catch an OOMError and re-run the same task on a bigger box, in a try block you wrote. The runtime knows what compute exists and can provision more mid-run.
A retry that starts from the top is not recovery, it is the same run again with the same bill. Union records what finished as the run happens, outside the node doing the work. That record is what you resume from, fork from, or run again with new code.
A pipeline’s output should not be a path you need to remember. Artifacts are typed, versioned values that persist past the run that made them, so a second workflow consumes the first one’s output without re-running it, and a new version can fire the next run by itself.
Every run keeps what it takes to reproduce a result: the code, the infrastructure it ran on, and the configuration applied to both. That is what makes a pipeline shareable rather than personal, and it is the same record an audit or a postmortem needs. For agent workloads it doubles as the trace of what the agent did.
Fanning out is plain async. asyncio.gather when you want the whole set at once, flyte.map.aio when you want bounded concurrency over a long list. If the cloud takes a spot node back mid-run, that branch retries and the rest keep going.
Workflow engines have handled a process dying for years. What is new is the demand that durability cover an OOM kill, a spot node taken back mid-run, & a flaky GPU. Here that is a try block that reaches the cluster: catch OOMError, raise the memory, run the same task again.
import flyteimport flyte.errors async def retry_with_memory( task_fn, *args, initial_memory: str = "250Mi", increment: str = "200Mi", max_memory: str = "4Gi", cpu: int = 1, **kwargs,): current_memory_mi = parse_memory(initial_memory) increment_mi = parse_memory(increment) max_memory_mi = parse_memory(max_memory) attempt = 1 while current_memory_mi <= max_memory_mi: mem_str = format_memory(current_memory_mi) print(f"Attempt {attempt}: running with memory: {mem_str}") try: result = await task_fn.override( resources=flyte.Resources(cpu=cpu, memory=mem_str) )(*args, **kwargs) print(f"Success with memory: {mem_str}") return result except flyte.errors.OOMError as e: print(f"OOMError with memory {mem_str}: {e}") if current_memory_mi + increment_mi > max_memory_mi: break current_memory_mi += increment_mi attempt += 1 raise RuntimeError( f"Task failed with OOM even after retrying up to " f"{format_memory(max_memory_mi)} across {attempt} attempts" )
except flyte.errors.OOMErrorAn OOM kill is an exception, not a dead run. So are TaskInterruptedError, TaskTimeoutError, and RetriesExhaustedError.
.override(resources=...)The handler changes the hardware for the next attempt. Infrastructure as context, with a signature.
A while loop and a try block. Write it yourself, share it as a helper, or let an agent write it against the same API.
Most of what makes a pipeline feel sluggish is startup: pulling an image, scheduling a pod, loading a model. An agent loop that does that on every step spends more time booting than thinking. Union keeps the expensive parts warm and pays the startup cost once.
Set reusable=flyte.ReusePolicy(replicas=4, concurrency=10) on a task environment and its containers stay up. The model loads once per process and every follow-on call shares it — a per-step cold start becomes a function call.
Once an action pays to spin up a pool or a Ray cluster, later actions land on that same cluster and its footprint is counted once, however many run inside it. Warm work rides its own lane, so a burst can’t crowd out the queue.
The image is declared in code next to the task that needs it, and the remote builder rebuilds it when that declaration changes. No docker build to run, and none to wait on.
When a replica does start cold, stream_model=True streams weights from your object store straight into GPU memory instead of staging them on local disk. Endpoints scale to zero between bursts without the restart being a coffee break.
Autonomous driving, cancer therapy, and agentic research at catalog scale. Long runs on infrastructure large enough to fail constantly.
Union is built on Flyte, the open-source AI runtime we create and maintain under the Linux Foundation AI & Data.
companies using Flyte today
Flyte SDK downloads
Bring a workload you have lost before, and watch what happens when the machine under it goes away.
Try for free