DEV Community

Yu Ting Chen
Yu Ting Chen

Posted on Edited on AI-assisted

When Celery Chains Get Hard to Maintain: Introducing CeleryFlow

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:

  1. Composing chains was painful. When place_order_flow and refund_flow needed the same opening sequence, you'd either duplicate the chain construction or wrap it in helper functions that ballooned over time.
  2. Conditionals were inline. "Only run step X when payload field Y looks like Z" became if statements inside the task body, mixing business logic with flow control.

Two-column comparison. Left, raw chain(): PlaceOrder and Refund each repeat the validate and charge steps (highlighted as duplicated), and the loyalty step carries an inline 'if returning' check. Right, CeleryFlow: a single shared checkout sub-flow (validate then charge) that both PlaceOrder and Refund point to, and the loyalty step carries a declarative condition instead of an inline if.

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" }
Enter fullscreen mode Exit fullscreen mode

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]"
Enter fullscreen mode Exit fullscreen mode

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.

Quadrant chart, y-axis durable execution low to high, x-axis separate service / infra weight low to high. Celery + chain(), CeleryFlow and arq/TaskIQ sit bottom-left; Prefect, Dagster and Airflow sit mid-right; Temporal and Step Functions sit top-right.

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.

Top-down decision tree. Multi-step background jobs? If the workflow must resume automatically after a crash, use durable execution (Temporal / Hatchet / Restate). Else if it is really a data pipeline, use Prefect / Dagster. Else if you are not already running Celery, use a plain task queue (Celery / arq / TaskIQ). Else if you have 5-50 multi-step flows with conditions or shared sub-flows, use CeleryFlow; otherwise plain Celery chain() / group().

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.