Temporal with Python: Durable Execution for Reliable Workflows
Temporal is a platform for durable execution: it lets you write long-running, stateful business logic as ordinary code that survives process crashes, deployments, and infrastructure failures. Instead of stitching together queues, cron jobs, and a database to track "where did this order get to," you write a function and Temporal guarantees it runs to completion exactly as written, even if the machine running it dies halfway through.
What Durable Execution Means
Most workflow tools you may have seen — Airflow, Prefect, Dagster — are schedulers for data pipelines. They are excellent at running a DAG of tasks on a cadence and showing you the results. Celery is a task queue for offloading background jobs. Temporal solves a different problem: keeping a single long-lived process correct across failures.
The core idea is the event-sourced replay model:
- A Workflow is your business logic written as deterministic code. Temporal does not keep the workflow's variables in memory forever. Instead, every meaningful step (an activity result, a timer firing, a signal arriving) is appended to a durable event history.
- If the worker process crashes, Temporal starts the workflow again on another worker and replays the event history. Each line of your code re-executes, but instead of calling activities again, the SDK feeds back the recorded results. When replay catches up to where the crash happened, execution continues normally.
- Because of replay, workflow code must be deterministic: given the same history, it must take the same path every time. Anything non-deterministic (network calls, random numbers, reading the clock, querying a database) must happen inside an Activity.
An Activity is a plain function for side effects and non-deterministic work. Activities are not replayed; their results are recorded once and reused. They are the place where your code talks to the outside world.
This split is what makes the magic work: the workflow is the durable, replayable "brain," and activities are the disposable "hands" that touch external systems.
Core Architecture
A Temporal deployment has a few moving parts:
- Temporal Server (Cluster): the backend that stores event histories, schedules tasks, and enforces timeouts and retries. It is the source of truth, backed by a database (PostgreSQL, MySQL, or Cassandra).
- Task Queues: named queues the server uses to hand work to your code. Workers poll a task queue; clients and the server route workflow and activity tasks onto it.
- Workers: processes you run that host your workflow and activity code. A worker polls a task queue, executes tasks, and reports results back to the server. Your code lives here, not on the server.
- Client: how application code starts workflows, sends signals, and queries state.
- Web UI: a dashboard to inspect running and completed workflows, their event histories, inputs, outputs, and failures.
The server never runs your code. It only orchestrates. This separation means you can deploy new worker versions, scale workers horizontally, and the server keeps the histories safe.
Getting a Development Server
For local development, the Temporal CLI ships a self-contained dev server with an in-memory database and the Web UI.
# Install the Temporal CLI (macOS / Linux)
curl -sSf https://temporal.download/cli.sh | sh
Or with Homebrew
brew install temporal
Start a local dev server with the Web UI on http://localhost:8233
temporal server start-dev
The dev server listens for SDK connections on localhost:7233 and serves the Web UI on localhost:8233. It resets state on restart, which is fine for development. For production you run a real cluster or use Temporal Cloud (covered later).
Install the Python SDK:
pip install temporalio
A Coherent Example: Order Fulfillment
Throughout this tutorial we build one workflow: processing a customer order. The steps are charge the payment, reserve inventory, ship the package, and notify the customer. Each step can fail and should be retried. We will later add a durable timer, a signal to cancel, and a query to check status.
Defining Activities
Activities hold all the side effects. Each is a normal (optionally async) function decorated with @activity.defn.
# activities.py
import asyncio
from dataclasses import dataclass
from temporalio import activity
@dataclass
class OrderInput:
orderid: str
customeremail: str
amountcents: int
@activity.defn
async def chargepayment(order: OrderInput) -> str:
activity.logger.info(f"Charging {order.amountcents} for {order.orderid}")
# Real code would call a payment gateway here.
await asyncio.sleep(0.2)
return f"charge{order.orderid}"
@activity.defn
async def reserveinventory(orderid: str) -> bool:
await asyncio.sleep(0.2)
return True
@activity.defn
async def shippackage(orderid: str) -> str:
await asyncio.sleep(0.5)
return f"tracking{orderid}"
@activity.defn
async def sendnotification(email: str, message: str) -> None:
activity.logger.info(f"Email to {email}: {message}")
Activities can fail and Temporal will retry them according to a policy. Because their results are durably recorded, a retried or replayed workflow never charges the same card twice for the same recorded step.
Defining the Workflow
The workflow orchestrates activities. It is decorated with @workflow.defn, and the entry point with @workflow.run.
# workflows.py
from datetime import timedelta
from temporalio import workflow
from temporalio.common import RetryPolicy
Import activity stubs through the sandbox-safe pass-through.
with workflow.unsafe.importspassedthrough():
from activities import (
OrderInput,
chargepayment,
reserveinventory,
shippackage,
sendnotification,
)
@workflow.defn
class OrderWorkflow:
def init(self) -> None:
self.status = "started"
@workflow.run
async def run(self, order: OrderInput) -> str:
retry = RetryPolicy(
initialinterval=timedelta(seconds=1),
backoffcoefficient=2.0,
maximuminterval=timedelta(seconds=30),
maximumattempts=5,
)
self.status = "charging"
chargeid = await workflow.executeactivity(
chargepayment,
order,
starttoclosetimeout=timedelta(seconds=10),
retrypolicy=retry,
)
self.status = "reserving"
await workflow.executeactivity(
reserveinventory,
order.orderid,
starttoclosetimeout=timedelta(seconds=10),
retrypolicy=retry,
)
self.status = "shipping"
tracking = await workflow.executeactivity(
shippackage,
order.orderid,
starttoclosetimeout=timedelta(seconds=30),
retrypolicy=retry,
)
await workflow.executeactivity(
sendnotification,
args=[order.customeremail, f"Order shipped: {tracking}"],
starttoclosetimeout=timedelta(seconds=10),
retrypolicy=retry,
)
self.status = "completed"
return tracking
Notice the workflow never calls requests, time.time(), or random directly. It only awaits activities and SDK primitives. That is the determinism constraint in practice.
The Determinism Constraint
When you need something non-deterministic inside workflow code, use the SDK's deterministic substitutes, which record their values into history:
# Inside a workflow — correct
now = workflow.now() # not datetime.now()
requestid = workflow.uuid4() # not uuid.uuid4()
delay = workflow.random().randint(1, 5) # not random.randint
Inside a workflow — WRONG, breaks replay
import datetime, uuid, random
now = datetime.datetime.now() # different value on replay -> non-determinism error
The Python SDK runs workflow code in a sandbox that blocks many of these mistakes, but you should still treat the rule as fundamental: no I/O, no clock, no randomness, no global mutable state in workflow code. Put it in an activity.
Running a Worker
A worker connects to the server, registers your workflows and activities, and polls a task queue.
# worker.py
import asyncio
from temporalio.client import Client
from temporalio.worker import Worker
from workflows import OrderWorkflow
from activities import (
chargepayment,
reserveinventory,
shippackage,
sendnotification,
)
async def main() -> None:
client = await Client.connect("localhost:7233")
worker = Worker(
client,
taskqueue="order-fulfillment",
workflows=[OrderWorkflow],
activities=[
chargepayment,
reserveinventory,
shippackage,
sendnotification,
],
)
await worker.run()
if name == "main":
asyncio.run(main())
Run it with python worker.py. The worker stays alive polling the order-fulfillment queue. Scaling throughput is as simple as starting more worker processes pointed at the same queue.
Starting a Workflow from a Client
Application code (an API handler, a script) starts workflows through a client.
# startorder.py
import asyncio
from temporalio.client import Client
from workflows import OrderWorkflow
from activities import OrderInput
async def main() -> None:
client = await Client.connect("localhost:7233")
order = OrderInput(
order
id="A-1001",
customeremail="customer@example.com",
amountcents=4999,
)
# executeworkflow starts and waits for the result.
result = await client.executeworkflow(
OrderWorkflow.run,
order,
id="order-A-1001", # the workflow ID
taskqueue="order-fulfillment",
)
print("Tracking:", result)
if name == "main":
asyncio.run(main())
Workflow IDs and Idempotency
The id is the unique identifier of a workflow execution. It is your idempotency key. If you try to start a workflow with an ID that is already running, Temporal rejects the duplicate by default. This means you can safely retry "start order A-1001" from an unreliable caller and never create two fulfillment processes for the same order. Use a business identifier (the order ID) rather than a random value.
To start without waiting, use startworkflow and get a handle:
handle = await client.startworkflow(
OrderWorkflow.run,
order,
id="order-A-1001",
task
queue="order-fulfillment",
)
result = await handle.result() # await later if needed
Retries and Timeouts
Temporal distinguishes several timeouts. The one you set most often is starttoclosetimeout: the maximum time a single activity attempt may run before it is considered failed and retried. The RetryPolicy controls how failures are retried, with exponential backoff.
from temporalio.common import RetryPolicy
from datetime import timedelta
retry = RetryPolicy(
initialinterval=timedelta(seconds=1),
backoffcoefficient=2.0,
maximuminterval=timedelta(minutes=1),
maximumattempts=10, # 0 means unlimited
nonretryableerrortypes=["InvalidCardError"],
)
By default activities retry indefinitely with backoff, so a transient payment-gateway outage simply resolves itself once the gateway recovers; the workflow waits patiently. Set maximumattempts to bound it.
Heartbeating Long Activities
For activities that run for a long time (processing a large file, calling a slow model), use heartbeats so the server knows the activity is still alive and can detect a stuck worker quickly. Pair heartbeats with a heartbeattimeout.
@activity.defn
async def processlargebatch(items: list[str]) -> int:
processed = 0
for i, item in enumerate(items):
# do work...
processed += 1
activity.heartbeat(i) # report progress; enables fast failure detection
return processed
await workflow.executeactivity(
process
largebatch,
items,
start
toclosetimeout=timedelta(minutes=30),
heartbeattimeout=timedelta(seconds=30),
)
If the worker crashes, the missed heartbeat lets the server reschedule the activity quickly instead of waiting out the full starttoclosetimeout. On retry, the activity can read activity.info().heartbeatdetails to resume from where it left off.
Durable Timers
A workflow can sleep for arbitrarily long — seconds, days, or months — without holding any resources. workflow.sleep creates a durable timer stored in history. The worker can shut down during the sleep; when the timer fires, Temporal wakes the workflow on any available worker.
@workflow.run
async def run(self, order: OrderInput) -> str:
# ... ship the package ...
# Wait 7 days, then ask for a review. No process needs to stay running.
await workflow.sleep(timedelta(days=7))
await workflow.executeactivity(
sendnotification,
args=[order.customeremail, "How was your order?"],
starttoclosetimeout=timedelta(seconds=10),
)
return tracking
This is fundamentally different from cron. Cron fires on a clock and you must look up state in a database to know what to do. A durable timer is part of one continuous workflow that already holds all its context in local variables. There is no external state to reconcile.
Signals, Queries, and Updates
Running workflows are not black boxes. You can interact with them.
A Signal sends data into a running workflow asynchronously. It is fire-and-forget and can change workflow state. Here we let a customer cancel before shipping.
@workflow.defn
class OrderWorkflow:
def init(self) -> None:
self.status = "started"
self.cancelled = False
@workflow.signal
def cancelorder(self) -> None:
self.cancelled = True
@workflow.query
def status(self) -> str:
return self.status
@workflow.run
async def run(self, order: OrderInput) -> str:
self.status = "charging"
await workflow.executeactivity(
chargepayment, order,
starttoclosetimeout=timedelta(seconds=10),
)
# Wait up to an hour for a possible cancel before shipping.
try:
await workflow.waitcondition(
lambda: self.cancelled,
timeout=timedelta(hours=1),
)
except TimeoutError:
pass
if self.cancelled:
self.status = "cancelled"
await workflow.executeactivity(
sendnotification,
args=[order.customeremail, "Your order was cancelled."],
starttoclosetimeout=timedelta(seconds=10),
)
return "cancelled"
self.status = "shipping"
tracking = await workflow.executeactivity(
shippackage, order.orderid,
starttoclosetimeout=timedelta(seconds=30),
)
self.status = "completed"
return tracking
A Query reads workflow state synchronously without changing it. Queries must not mutate state or call activities. An Update is the newer, validated, request/response interaction that can both change state and return a result, with optional validation before acceptance.
From the client:
handle = client.getworkflowhandle("order-A-1001")
Query current status
current = await handle.query(OrderWorkflow.status)
print(current)
Signal a cancellation
await handle.signal(OrderWorkflow.cancelorder)
Child Workflows and continueasnew
A workflow can start child workflows to decompose large problems or to give a sub-process its own ID, retry policy, and lifecycle.
tracking = await workflow.executechildworkflow(
ShippingWorkflow.run,
order.order
id,
id=f"shipping-{order.orderid}",
)
The event history grows with every step. A workflow that loops forever (a subscription that bills monthly, an agent that runs many iterations) would eventually accumulate an enormous history, which slows replay. The fix is continueasnew: it ends the current execution and atomically starts a fresh one with the same ID and a clean history, carrying forward only the state you pass.
@workflow.run
async def run(self, state: int) -> None:
for in range(1000):
await workflow.executeactivity(
doperiodicwork,
starttoclosetimeout=timedelta(seconds=30),
)
await workflow.sleep(timedelta(days=30))
state += 1
# Truncate history and continue with carried-over state.
workflow.continueasnew(state)
Error Handling
Activity exceptions surface in the workflow as ActivityError wrapping an ApplicationError. To mark an error as permanent so Temporal does not retry it, raise a non-retryable ApplicationError from the activity.
from temporalio.exceptions import ApplicationError
@activity.defn
async def chargepayment(order: OrderInput) -> str:
if order.amountcents <= 0:
# A bad request will never succeed on retry.
raise ApplicationError("Invalid amount", nonretryable=True)
...
In the workflow you can catch and compensate, which is how you build a saga (undo previous steps when a later one fails):
from temporalio.exceptions import ActivityError
try:
await workflow.executeactivity(reserveinventory, order.orderid, ...)
except ActivityError:
# Compensate: refund the earlier charge.
await workflow.executeactivity(refundpayment, chargeid, ...)
raise
Versioning and Patching
Because workflows replay old histories, changing workflow code can break in-flight executions: the new code might take a different path than the recorded history expects. For safe changes, use workflow.patched, which records which branch a given execution took.
if workflow.patched("use-express-shipping"):
tracking = await workflow.executeactivity(expressship, ..., )
else:
tracking = await workflow.executeactivity(shippackage, ..., )
Old executions replay the original branch; new executions take the new one. Once all old workflows have completed, you call workflow.deprecatepatch("use-express-shipping") and later remove the conditional. For many teams, running workers with Worker Versioning (pinning a build ID) is a complementary strategy.
Testing with the Time-Skipping Environment
Temporal ships a test environment that runs an in-memory server and can skip time, so a workflow with a seven-day timer completes in milliseconds in your test suite.
# testorder.py
import uuid
import pytest
from temporalio.testing import WorkflowEnvironment
from temporalio.worker import Worker
from workflows import OrderWorkflow
from activities import (
OrderInput, charge
payment, reserveinventory,
ship
package, sendnotification,
)
@pytest.mark.asyncio
async def test
ordercompletes():
async with await WorkflowEnvironment.start
timeskipping() as env:
async with Worker(
env.client,
task
queue="test-queue",
workflows=[OrderWorkflow],
activities=[chargepayment, reserveinventory,
shippackage, sendnotification],
):
order = OrderInput("A-1", "a@b.com", 4999)
result = await env.client.executeworkflow(
OrderWorkflow.run,
order,
id=f"test-{uuid.uuid4()}",
taskqueue="test-queue",
)
assert result.startswith("tracking")
You can also mock activities by registering substitute functions on the test worker, letting you assert workflow logic without touching real external systems.
A Note on AI and Agent Workflows
Durable execution fits modern AI pipelines well. A multi-step LLM agent — retrieve context, call a model, run a tool, evaluate, maybe loop, then call the model again — is exactly the kind of long-running, failure-prone, stateful process Temporal is built for. Model calls and tool invocations become activities with their own retries and timeouts, so a rate-limit error or a flaky tool does not lose the whole run. The orchestration logic (which step comes next, when to stop, how many iterations) lives in a deterministic workflow, and continueasnew keeps a long agent loop's history bounded. A human-in-the-loop approval becomes a signal the workflow waits on with a durable timer as a deadline.
The same shape covers business sagas: an order, a loan application, or an onboarding flow that spans days, calls many services, and must compensate cleanly on failure.
Temporal Cloud vs Self-Hosting
You can run the Temporal Server yourself on Kubernetes or VMs, backed by your own PostgreSQL or Cassandra. This gives full control but means you operate the cluster, its database, scaling, and upgrades. Temporal Cloud is the managed offering: you run only your workers and connect them to a hosted namespace over mTLS, while Temporal operates the server, storage, and scaling. A common path is to start on the local dev server, validate on a small self-hosted cluster or Cloud namespace, and choose based on whether operating the backend is worth the control it gives you.
Best Practices and Common Pitfalls
- Keep workflows deterministic. No direct clock, randomness, network, file, or database access in workflow code. Use
workflow.now,workflow.uuid4,workflow.random, and activities. Determinism violations are the most common production bug. - Put all side effects in activities. If it touches the outside world or could differ between runs, it belongs in an activity.
- Use meaningful workflow IDs. Tie the ID to a business entity for natural idempotency and easy lookup in the Web UI.
- Set explicit timeouts. Always set
starttoclosetimeouton activities; addheartbeattimeoutfor long ones and heartbeat from inside them. - Mark permanent failures non-retryable. A malformed input or a 4xx response should raise
nonretryable=Trueso you do not retry forever. - Bound long histories with continueasnew. Any loop or long-lived workflow should periodically continue-as-new.
- Version risky changes with patching. Never reshape the order of activities in a deployed workflow without
workflow.patchedwhile old executions are still running. - Do not block the workflow thread. Avoid
time.sleepand synchronous blocking calls; useworkflow.sleepand async activities. - Test with time skipping. Cover timer-heavy and signal-driven paths in the test environment rather than waiting in real time.
Conclusion
Temporal turns fragile, multi-step processes into ordinary code that simply does not lose progress. By splitting logic into deterministic workflows and side-effecting activities, and by reconstructing state through event-history replay, it removes the usual scaffolding of queues, status tables, and reconciliation jobs.
Key takeaways:
- Durable execution means your workflow code survives crashes and restarts via event-sourced replay.
- Workflows must be deterministic; all I/O and non-determinism go in activities.
- The server orchestrates and stores history; workers run your code; task queues connect them.
- Retries, timeouts, heartbeats, and durable timers are built-in, not bolted on.
- Signals, queries, and updates let you interact with live workflows; child workflows and
continueasnewkeep large processes manageable. - It suits long-running business sagas and multi-step AI/agent pipelines far better than a data-pipeline scheduler or a bare task queue.