If you run Celery in production and your multi-step jobs have turned
into a pile of chain() calls nobody wants to touch, this post is
about that specific problem.
I ran into this problem at a previous job and built a small internal library to solve it. Earlier this year, I released an open-source version as CeleryFlow. It lets you describe Celery workflows in YAML instead of assembling them in Python.
The honest disclaimer first: it's not the right tool for most new projects, and it is not a durable execution engine. If you're starting fresh in 2026, there are better places to look, and I'll point at them near the end. What follows is written for people already running
Celery.
When you outgrow chain()
At that job, we had dozens of multi-step workflows in an e-commerce platform: order placement, refund, plan upgrade, subscription renewal, etc. Each was 5-10 steps. They shared common sub-flows ("validate order", "charge card", "send receipt"). Conditions like "only run the loyalty-points step if buyer is a returning customer" appeared all over.
Celery's chain(), group(), and chord() primitives are powerful,
and none of this was impossible. But this approach had two specific pain points:
-
Composing chains was painful. When
place_order_flowandrefund_flowneeded the same opening sequence, you'd either duplicate the chain construction or wrap it in helper functions that ballooned over time. -
Conditionals were inline. "Only run step X when payload field
Y looks like Z" became
ifstatements inside the task body, mixing business logic with flow control.
Same two flows, two approaches. On the left, the steps are duplicated across flows, and the condition is implemented as an inline if statement. On the right, the sub-flow is defined once, and the condition is defined in config.
The natural response was to pull the workflow structure out of the code and into config files. We did, and it stuck. CeleryFlow brings that design to open source.
What it looks like
work-flows:
- name: checkout
tasks:
- order.validate
- order.charge
main-flows:
- name: PlaceOrder
flows:
- flow: checkout
- task: order.send_receipt
- task: order.award_loyalty_points
condition:
customer_type: { $eq: "returning" }
That's the whole workflow definition, but not the whole setup. To wire it in, use the CeleryFlow app class instead of Celery, give participating tasks base=EventTask, and register the file with FlowBuilder.from_yaml(). Once registered, the flow runs on your existing Celery workers and broker. There is no new worker model or orchestration server, and your existing retry behavior, routing, Flower setup, and monitoring continue to work.
One important caveat: in v0.2.0, a condition that doesn't match raises ConditionFailed before the task is queued. There is no built-in quiet skip; unless the exception is handled explicitly, the chain fails. A missing field counts as passing, so conditions only gate fields that are present in the payload.
pip install "celeryflow[yaml]"
The yaml extra pulls in PyYAML, which FlowBuilder.from_yaml()
needs. Plain pip install celeryflow works too if you'd rather keep
your flow definitions as Python dicts or JSON.
Quickstart and full docs: https://github.com/ChenYuTingJerry/CeleryFlow
When you don't need it
You don't need any workflow framework if:
- You only have single tasks or short chains with one or two follow-ups, and you're happy using
Task.s() | Task.s()for the chains. - You don't have many of them.
- The chains rarely change shape.
For a handful of simple workflows, just use the Celery primitives. Adding a layer on top is overhead you don't need.
What it isn't
The most important question before choosing an approach is:
If a worker process dies in the middle of a multi-step workflow,
what happens?
With Celery, and therefore with CeleryFlow, individual tasks may be redelivered after certain failures, depending on acks_late, the broker, and how the worker failed. A task may also run more than once. What Celery doesn't provide is durable workflow state from which the whole flow can resume.
If that's not acceptable, look at durable execution platforms:
Temporal, Hatchet, Restate, AWS Step Functions. They persist workflow
progress, so execution can recover after a failure and wait between
steps. Adopting one means moving to a different execution and operational model, whether self-hosted or managed. CeleryFlow has no equivalent persistence or durable timers, so it is a poor fit for flows that wait for hours or days between steps. It is also Python-only, while most of those platforms support several languages.
Durable workflow recovery does not mean that every side effect happens exactly once. Individual operations may still be retried, so anything involving money or an external API needs to be idempotent. The exact guarantees vary by platform.
And if your "workflow" is really a data pipeline (ETL, batch ML
training, scheduled report generation), Prefect and Dagster are built
for that, with first-class scheduling, observability UIs, and
integrations with warehouses and dbt. CeleryFlow can run scheduled
tasks via Celery Beat, but that's not what it's for.
Where it sits
My first-hand experience here is with Celery and CeleryFlow. I've also operated Airflow at the infrastructure level. For the other tools, I'm relying on their intended design rather than direct use.
| Tool | Type | When to use it |
|---|---|---|
| Celery | Task queue | You need workers that run individual tasks. Mature, battle-tested. |
Celery + chain()/group() |
Manual workflow | You have Celery and a handful of multi-step jobs. Works, gets messy at scale. |
| CeleryFlow | Config-driven workflow on Celery | You already use Celery, want declarative workflows, don't want a separate service. |
| Temporal | Durable execution platform | You need workflows that survive worker crashes, run for days, and resume from where they stopped. |
| Prefect | Modern data orchestration | Data pipelines, ML workflows, want a UI, want hybrid cloud / on-prem. |
| Dagster | Asset-oriented orchestration | Data engineering teams who think in terms of "data assets" not "tasks." |
| Airflow | Classic DAG scheduler | You're at a company that already runs Airflow, or your team is comfortable with it. |
| arq / TaskIQ / Procrastinate | Async-native task queues | Brand-new project, async-first, don't need Celery's surface area. |
Two questions capture most of the trade-off: do you need durable execution, and how much additional infrastructure are you willing to run?
So is it for you?
I'll be specific. CeleryFlow is a good fit if:
- You have a Python backend that already runs Celery in production.
- You have between 5 and 50 multi-step workflows.
- Each workflow is seconds to minutes long, not days.
- Workflows have conditions that gate a step on the payload, or share common sub-flows.
- You don't want a separate service for workflow orchestration.
- Re-triggering the whole flow from the upstream event is an acceptable recovery strategy.
The whole article as one picture. CeleryFlow is the leaf at the end of a fairly narrow path, and that's on purpose.
What I'd like feedback on
I'm trying to understand whether this design still fits real-world needs in 2026 or whether teams have moved on. If you've used Celery at any scale, I'd love to hear your perspective:
- Have you ever written something like CeleryFlow? Why or why not?
- Did you switch from Celery to Temporal / Prefect / something else? What pushed you?
- If you tried CeleryFlow and didn't keep it, what made you bounce?
GitHub issues and comments here on DEV are both welcome.
I used AI to help draft and edit the English. I checked the technical claims against the CeleryFlow code.



Top comments (1)
Some comments may only be visible to logged-in visitors. Sign in to view all comments.