Temporal Tutorial: Durable Execution for Reliable Workflows

# Temporal dengan Python: Durable Execution untuk Workflow yang Andal Temporal adalah platform untuk durable execution: ia memungkinkan Anda menulis logika bisnis yang berjalan lama dan bersifat stat...

By Ruby Abdullah · · tutorial
TemporalDurable ExecutionWorkflowOrchestrationDistributed SystemsPython

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(

orderid="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",

taskqueue="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(

processlargebatch,

items,

starttoclosetimeout=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.orderid,

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, chargepayment, reserveinventory,

shippackage, sendnotification,

)

@pytest.mark.asyncio

async def testordercompletes():

async with await WorkflowEnvironment.starttimeskipping() as env:

async with Worker(

env.client,

taskqueue="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 starttoclosetimeout on activities; add heartbeattimeout for long ones and heartbeat from inside them.
  • Mark permanent failures non-retryable. A malformed input or a 4xx response should raise nonretryable=True so 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.patched while old executions are still running.
  • Do not block the workflow thread. Avoid time.sleep and synchronous blocking calls; use workflow.sleep and 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 continueasnew keep 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.

Related Articles

ComfyUI Tutorial: Node-Based Workflows for Stable Diffusion

ComfyUI: Workflow Berbasis Node untuk Stable Diffusion ComfyUI adalah lingkungan grafis berbasis node untuk menjalankan ...

ZenML: Modular and Cloud-Agnostic MLOps Pipeline Framework

ZenML: Framework Pipeline MLOps yang Modular dan Cloud-Agnostic Pendahuluan Membangun model machine learning yang akurat...

AWS Step Functions for ML Tutorial: ML Workflow Orchestration

Tutorial Lengkap AWS Step Functions untuk ML: Orkestrasi ML Workflows AWS Step Functions menyediakan orkestrasi workflow...

Complete Apache Airflow Tutorial: Workflow Orchestration for Data Pipelines

Tutorial Lengkap Apache Airflow: Workflow Orchestration untuk Data Pipelines Apache Airflow adalah platform open-source ...