Migrate a Celery Task Queue to a Temporal Workflow
Celery is a distributed task queue that runs background jobs by pushing messages through a broker (such as Redis or RabbitMQ) to a pool of worker processes. It is a popular choice for sending emails, processing uploads, and running scheduled jobs. As applications grow, however, teams often want stronger guarantees than a broker provides: automatic recovery after a crash, a durable record of every job's progress, and built-in retries.
Temporal provides those guarantees. Instead of enqueueing a message and waiting for a worker to finish it, you run a Workflow whose entire state is persisted by the Temporal service. If a Worker crashes mid-job, another Worker resumes exactly where it left off.
In this guide, you will migrate a Celery application to a Temporal Workflow one piece at a time. You will convert a Celery task into a Temporal Activity, orchestrate it with a Workflow, run a Worker to process it, and replace your .delay() calls, task retries, Celery Beat schedules, and (optionally) Canvas workflows with their Temporal equivalents. By the end, you will have a working Temporal Application that reproduces the behavior of your Celery app.
How Celery concepts map to Temporal
Before you start, it helps to know which Temporal building block replaces each Celery concept. You will implement each row of this table in the steps that follow.
| Celery | Temporal | Purpose |
|---|---|---|
Task (@app.task) | Activity (@activity.defn) | A single unit of non-deterministic work (I/O, API calls) |
| (no direct equivalent) | Workflow (@workflow.defn) | Durable orchestration that calls Activities |
Worker (celery worker) | Worker (temporalio.worker.Worker) | Process that executes your code |
| Broker + result backend | Temporal service | Stores queue state, results, and history |
task.delay(...) | client.start_workflow(...) | Kick off work |
AsyncResult.get() | handle.result() | Retrieve the return value |
max_retries / self.retry() | RetryPolicy | Automatic retries |
| Celery Beat | Temporal Schedule | Recurring jobs |
Canvas (chain, group, chord) | Activity calls inside a Workflow | Multi-step pipelines |
Prerequisites
Before you begin, you will need the following:
- Python 3.14 or higher installed on your machine.
- An existing Celery application you want to migrate, or the sample Celery task shown in Step 4 if you are following along from scratch.
- Familiarity with running Python scripts from the command line.
Step 1 — Set up your project directory
In this step, you will create a project layout that keeps Workflow code and Activity code in separate files. Temporal reloads Workflow files frequently to protect determinism, so keeping them small improves Worker performance.
Create a new project directory and move into it:
mkdir temporal-migration && cd temporal-migration
Create the two package directories that will hold your code:
mkdir workflows activities
Your project will grow into the following structure as you work through the tutorial:
temporal-migration/
├── activities/ # Activity functions (your former Celery tasks)
├── workflows/ # Workflow classes (the orchestration layer)
├── worker.py # Runs the Worker
└── starter.py # Triggers Workflows
With the directory in place, you can install the tools you need.
Step 2 — Install the Temporal SDK and CLI
In this step, you will install the Python SDK your code depends on and the Temporal CLI you will use to run a local server.
Install the Temporal Python SDK with pip:
pip install temporalio
Next, install the Temporal CLI. On macOS or Linux with Homebrew, run:
brew install temporal
If you are not using Homebrew, download the binary for your platform from the Temporal CLI install guide and add it to your PATH.
Verify the CLI is available:
temporal --version
You will see the installed version printed to your terminal. With the tools installed, you can start a local Temporal service.
Step 3 — Start the Temporal development server
In Celery, work flows through a broker such as Redis. In Temporal, work flows through the Temporal service, which also stores each Workflow's durable history. In this step, you will start a local development server that stands in for that service.
Start the development server:
temporal server start-dev
You will see output confirming the server is running, including two addresses:
[secondary_label Output]
Server: localhost:7233
UI: http://localhost:8233
Your application code will connect to localhost:7233. The Web UI at http://localhost:8233 lets you inspect every Workflow, its inputs and outputs, and its complete event history — the equivalent of a more-detailed Flower dashboard.
Leave this process running and open a new terminal for the remaining steps.
Step 4 — Convert a Celery task into an Activity
In this step, you will take a Celery task and rewrite it as a Temporal Activity. An Activity holds the same non-deterministic work your task did, such as network calls, database writes, and file I/O, and Temporal takes responsibility for retrying it when it fails.
If you don’t have a Celery task of your own, consider a typical Celery task that sends a welcome email. Its tasks.py might look like this:
# Celery version — tasks.py
from celery import Celery
app = Celery("myapp", broker="redis://localhost:6379/0")
@app.task(bind=True, max_retries=5, default_retry_delay=10)
def send_welcome_email(self, user_id):
try:
user = get_user(user_id)
deliver_email(user.email, "Welcome!")
return f"sent to {user.email}"
except TransientError as exc:
raise self.retry(exc=exc)
Create activities/email.py and add the Temporal equivalent:
# activities/email.py
from dataclasses import dataclass
from temporalio import activity
@dataclass
class User:
user_id: int
email: str
@dataclass
class WelcomeEmailInput:
user_id: int
# --- Mock helpers ---------------------------------------------------------
# Stand-ins for your real user lookup and email delivery. Replace these with
# your database query and email provider when you adapt the guide.
def get_user(user_id: int) -> User:
return User(user_id=user_id, email=f"user{user_id}@example.com")
def deliver_email(address: str, subject: str) -> None:
print(f"Delivering '{subject}' to {address}")
# --------------------------------------------------------------------------
@activity.defn
def send_welcome_email(input: WelcomeEmailInput) -> str:
user = get_user(input.user_id)
deliver_email(user.email, "Welcome!")
return f"sent to {user.email}"
Two changes are worth noting. First, the retry boilerplate is gone: you no longer catch TransientError or call self.retry(), because Temporal retries a failed Activity automatically. You will configure how it retries in Step 8. Second, the Activity takes a single dataclass argument instead of positional parameters. Passing one structured argument is the recommended Temporal pattern, because it lets you add fields later without breaking running Workflows.
This Activity is a plain synchronous function, which is the safest default. Now you need something to invoke it.
Step 5 — Orchestrate the Activity with a Workflow
Celery has no equivalent of a Workflow — tasks are called directly. In Temporal, a Workflow is the durable coordinator that decides which Activities run and in what order. In this step, you will write a Workflow that calls the Activity from the previous step.
Create workflows/onboarding.py:
# workflows/onboarding.py
from datetime import timedelta
from temporalio import workflow
with workflow.unsafe.imports_passed_through():
from activities.email import send_welcome_email, WelcomeEmailInput
@workflow.defn
class OnboardingWorkflow:
@workflow.run
async def run(self, user_id: int) -> str:
return await workflow.execute_activity(
send_welcome_email,
WelcomeEmailInput(user_id),
start_to_close_timeout=timedelta(seconds=30),
)
Three details matter here. You import the Activity inside workflow.unsafe.imports_passed_through() so the Workflow sandbox does not reload it. The run method is decorated with @workflow.run and must be async. And every Activity call requires a timeout: start_to_close_timeout sets the maximum time a single attempt may run, replacing Celery's task_time_limit.
With the Workflow defined, you need a Worker to execute both it and the Activity.
Step 6 — Run a Worker to process Tasks
Just as celery worker pulls jobs from a broker, a Temporal Worker polls a Task Queue for work. In this step, you will register your Workflow and Activity with a Worker and start it.
Create worker.py in the project root:
# worker.py
import asyncio
import concurrent.futures
from temporalio.client import Client
from temporalio.worker import Worker
from activities.email import send_welcome_email
from workflows.onboarding import OnboardingWorkflow
async def main():
client = await Client.connect("localhost:7233")
with concurrent.futures.ThreadPoolExecutor(max_workers=100) as activity_executor:
worker = Worker(
client,
task_queue="onboarding",
workflows=[OnboardingWorkflow],
activities=[send_welcome_email],
activity_executor=activity_executor,
)
await worker.run()
if __name__ == "__main__":
asyncio.run(main())
The task_queue name is the routing key that ties your Worker, Workflow, and trigger code together, similar to a Celery queue name. Because your Activity is synchronous, you pass a ThreadPoolExecutor as the activity_executor; the max_workers value controls how many Activities run concurrently, much like Celery's --concurrency flag.
Start the Worker:
python worker.py
The Worker begins polling the onboarding Task Queue and waits for work. Leave it running and open another terminal to trigger it.
Step 7 — Trigger Workflows in place of .delay()
In Celery, you enqueue a job by calling send_welcome_email.delay(42). In Temporal, you start a Workflow through a client. In this step, you will replace your enqueue calls with start_workflow.
Create starter.py in the project root:
# starter.py
import asyncio
import uuid
from temporalio.client import Client
from workflows.onboarding import OnboardingWorkflow
async def main():
client = await Client.connect("localhost:7233")
handle = await client.start_workflow(
OnboardingWorkflow.run,
42,
id=f"onboarding-{uuid.uuid4()}",
task_queue="onboarding",
)
print(f"Started workflow {handle.id}")
result = await handle.result()
print(f"Result: {result}")
if __name__ == "__main__":
asyncio.run(main())
Run it:
python starter.py
You will see output confirming the Workflow ran to completion:
[secondary_label Output]
Started workflow onboarding-3f9a...
Result: sent to user@example.com
Note how the Celery patterns translate. A fire-and-forget .delay() corresponds to start_workflow, which returns a handle immediately. Retrieving the return value with AsyncResult.get() corresponds to handle.result(). The id you provide is a business identifier you choose (an order number, a user ID); Temporal uses it to guarantee that the same Workflow is never started twice, which is a built-in form of deduplication.
To schedule a Workflow to begin after a delay, as with Celery's apply_async(countdown=...), pass start_delay=timedelta(...) to start_workflow.
Step 8 — Migrate task retries to a Retry Policy
In Celery, retries are your responsibility: you set max_retries and call self.retry() inside the task. In Temporal, retries are automatic and declarative. In this step, you will restore your task's retry behavior by attaching a retry policy to the Activity call.
By default, Temporal retries a failed Activity indefinitely with exponential backoff. To reproduce the Celery task's limit of five attempts, update the execute_activity call in workflows/onboarding.py:
# workflows/onboarding.py
from datetime import timedelta
from temporalio import workflow
from temporalio.common import RetryPolicy
with workflow.unsafe.imports_passed_through():
from activities.email import send_welcome_email, WelcomeEmailInput
@workflow.defn
class OnboardingWorkflow:
@workflow.run
async def run(self, user_id: int) -> str:
return await workflow.execute_activity(
send_welcome_email,
WelcomeEmailInput(user_id),
start_to_close_timeout=timedelta(seconds=30),
retry_policy=RetryPolicy(
maximum_attempts=5,
maximum_interval=timedelta(minutes=1),
non_retryable_error_types=["InvalidUserError"],
),
)
Here, maximum_attempts=5 mirrors Celery's max_retries, and maximum_interval caps the backoff between attempts. The non_retryable_error_types list names errors that should fail immediately without retrying — the equivalent of not calling self.retry() for a permanent failure. To raise such an error from an Activity, use ApplicationError with non_retryable=True:
# activities/email.py (excerpt)
from temporalio.exceptions import ApplicationError
@activity.defn
def send_welcome_email(input: WelcomeEmailInput) -> str:
user = get_user(input.user_id)
if user is None:
raise ApplicationError("No such user", type="InvalidUserError", non_retryable=True)
deliver_email(user.email, "Welcome!")
return f"sent to {user.email}"
Unless you have a specific reason to change them, leaving the other retry options at their defaults is recommended. With retries in place, you can move on to scheduled work.
Step 9 — Replace Celery Beat with a Temporal Schedule
If your Celery app runs periodic tasks through Celery Beat, you need a recurring trigger in Temporal. A Temporal Schedule starts a Workflow on a recurring basis, and unlike Beat it requires no separate long-running scheduler process. In this step, you will create a Schedule that runs a Workflow once a day.
First, create the Workflow the Schedule will run, along with a small Activity for it to call. Add activities/reports.py:
# activities/reports.py
from temporalio import activity
@activity.defn
def generate_daily_report() -> str:
# Replace with your real report logic (query metrics, render a file, etc.)
return "report generated"
Then add workflows/reports.py:
# workflows/reports.py
from datetime import timedelta
from temporalio import workflow
with workflow.unsafe.imports_passed_through():
from activities.reports import generate_daily_report
@workflow.defn
class DailyReportWorkflow:
@workflow.run
async def run(self) -> str:
return await workflow.execute_activity(
generate_daily_report,
start_to_close_timeout=timedelta(minutes=5),
)
A Schedule only starts a Workflow — a Worker still has to execute it. Register the new Workflow and Activity with the Worker you built in Step 6 by adding them to the lists in worker.py:
# worker.py (updated imports and registration)
from activities.email import send_welcome_email
from activities.reports import generate_daily_report
from workflows.onboarding import OnboardingWorkflow
from workflows.reports import DailyReportWorkflow
# ... inside main(), update the Worker: ...
worker = Worker(
client,
task_queue="onboarding",
workflows=[OnboardingWorkflow, DailyReportWorkflow],
activities=[send_welcome_email, generate_daily_report],
activity_executor=activity_executor,
)
Restart the Worker so it picks up the new registrations:
python worker.py
Create schedule.py in the project root:
# schedule.py
import asyncio
from datetime import timedelta
from temporalio.client import (
Client,
Schedule,
ScheduleActionStartWorkflow,
ScheduleSpec,
ScheduleIntervalSpec,
)
from workflows.reports import DailyReportWorkflow
async def main():
client = await Client.connect("localhost:7233")
await client.create_schedule(
"daily-report",
Schedule(
action=ScheduleActionStartWorkflow(
DailyReportWorkflow.run,
id="daily-report",
task_queue="onboarding",
),
spec=ScheduleSpec(
intervals=[ScheduleIntervalSpec(every=timedelta(days=1))],
),
),
)
print("Schedule created")
if __name__ == "__main__":
asyncio.run(main())
Run it once to register the Schedule:
python schedule.py
The ScheduleIntervalSpec shown here fires every 24 hours. For the calendar-style timing you may have expressed with Beat's crontab(...), the ScheduleSpec also accepts calendar and cron specifications; see the Temporal documentation for the full set of options. Once registered, a Schedule can be paused, resumed, or triggered on demand from the Web UI or through the client, replacing Beat's static configuration file.
Step 10 — (Optional) Translate Canvas workflows
Celery's Canvas primitives — chain, group, and chord — let you compose tasks into pipelines. In Temporal, this composition lives in ordinary Python inside a Workflow, which makes multi-step logic easier to read and debug. In this step, you will translate the common Canvas patterns.
A Celery chain runs tasks in sequence, passing each result to the next:
# Celery version
chain(fetch.s(url), parse.s(), store.s())()
In a Workflow, sequencing is successive await calls:
# workflows/pipeline.py (chain equivalent)
raw = await workflow.execute_activity(fetch, url, start_to_close_timeout=timedelta(minutes=1))
parsed = await workflow.execute_activity(parse, raw, start_to_close_timeout=timedelta(minutes=1))
result = await workflow.execute_activity(store, parsed, start_to_close_timeout=timedelta(minutes=1))
A Celery group runs tasks in parallel, and a chord runs a callback once a group finishes. Both map to asyncio.gather followed by an optional final Activity:
# workflows/pipeline.py (group + chord equivalent)
import asyncio
from datetime import timedelta
from temporalio import workflow
with workflow.unsafe.imports_passed_through():
from activities.files import process_file, summarize, FileInput, SummaryInput
@workflow.defn
class ImportWorkflow:
@workflow.run
async def run(self, file_ids: list[int]) -> str:
# group: process every file in parallel
results = await asyncio.gather(*[
workflow.execute_activity(
process_file,
FileInput(file_id),
start_to_close_timeout=timedelta(minutes=5),
)
for file_id in file_ids
])
# chord callback: run once all parallel work is done
return await workflow.execute_activity(
summarize,
SummaryInput(results),
start_to_close_timeout=timedelta(minutes=1),
)
Because this orchestration is plain Python, you can add conditionals, loops, and error handling around it without learning new Canvas syntax. For very large pipelines, Temporal offers child workflows so a Workflow can start other Workflows; see the Temporal documentation for that pattern.
Conclusion
In this guide, you migrated a Celery application to Temporal. You converted a Celery task into an Activity, wrapped it in a durable Workflow, ran a Worker to execute both, and replaced your .delay() calls, retries, Beat schedules, and Canvas pipelines with Temporal equivalents. Your jobs now survive worker crashes automatically, keep a complete history you can inspect in the Web UI, and retry on well-defined policies.
To continue, consider running your Workflows and Activities against production-style infrastructure and adding automated tests. Useful next topics include:
- Retrying and timing out work correctly, in the Temporal documentation.
- Self-hosting the Temporal service or using Temporal Cloud instead of the development server.
- The original Celery documentation for confirming the exact behavior of the tasks you are migrating.