Describe Celery workflows in YAML, JSON, or a Python dict, with optional per-step conditions.
Project description
CeleryFlow
CeleryFlow lets you describe a multi-step Celery workflow as a config
file (YAML, JSON, or a plain Python dict) instead of wiring chain(),
group(), and chord() by hand. Each step can have a condition, so
you can skip parts of the flow based on the payload at runtime.
I used a version of this internally for a few years to run order and billing pipelines. This is the cleaned-up open-source version with the business-specific stuff removed.
What you get
- Define flows in YAML, JSON, or a Python dict.
- Reuse named sub-flows inside larger flows.
- Skip steps with conditions like
{"sn": {"$eq": "1234"}}—$eq,$ne,$in,$gte, etc. You can add your own. - Optional structured logging on start / success / failure / retry.
- The flows are still regular Celery tasks. Retries, routing, and monitoring keep working the way you expect.
Install
pip install celeryflow # core
pip install "celeryflow[yaml]" # if you want YAML configs
pip install "celeryflow[redis]" # if you want Redis as the broker
Python 3.11 or newer. During development I use uv, but you don't need it to install the package.
Quickstart
from celeryflow import CeleryFlow, EventTask, FlowBuilder
app = CeleryFlow("my_worker")
app.conf.update(broker_url="redis://localhost:6379/0",
result_backend="redis://localhost:6379/0")
# 1. Write some tasks.
@app.task(name="order.validate", base=EventTask)
def validate(payload):
assert "sn" in payload
return payload
@app.task(name="order.charge", base=EventTask)
def charge(payload):
payload["charged"] = True
return payload
@app.task(name="order.notify", base=EventTask)
def notify(payload):
payload["notified"] = True
return payload
- Describe the flow. CeleryFlow accepts the definition in three formats — a Python dict, a YAML file, or a JSON file. Pick whichever fits your project.
Option A — Python dict, inline:
FlowBuilder.from_dict(app, {
"work-flows": [
{"name": "checkout", "tasks": ["order.validate", "order.charge"]},
],
"main-flows": [
{
"name": "PlaceOrder",
"flows": [
{"flow": "checkout"},
{"task": "order.notify"},
],
},
],
})
Option B — YAML file:
# flows.yaml
work-flows:
- name: checkout
tasks:
- order.validate
- order.charge
main-flows:
- name: PlaceOrder
flows:
- flow: checkout
- task: order.notify
FlowBuilder.from_yaml(app, "flows.yaml")
# Globs work too — multiple files get merged:
FlowBuilder.from_yaml(app, "config/flows/*.yaml")
Option C — JSON file:
{
"work-flows": [
{"name": "checkout", "tasks": ["order.validate", "order.charge"]}
],
"main-flows": [
{
"name": "PlaceOrder",
"flows": [
{"flow": "checkout"},
{"task": "order.notify"}
]
}
]
}
FlowBuilder.from_json(app, "flows.json")
All three produce the same PlaceOrder task. Then run it:
# 3. Run it.
result = app.execute_flow("PlaceOrder", {"sn": "1234"})
print(result.get())
# -> {"sn": "1234", "charged": True, "notified": True}
Conditional steps
FlowBuilder.from_dict(app, {
"main-flows": [
{
"name": "RefundFlow",
"flows": [
{"task": "order.validate"},
# Only charge when it's a new subscription worth >= 100.
{"task": "order.charge",
"condition": {"buy_action": {"$eq": "NEW_SUBSCRIPTION"},
"amount": {"$gte": 100}}},
],
},
],
})
When a condition doesn't match, apply_async raises
celeryflow.exceptions.ConditionFailed. Catch it in a link_error or
custom callback if you'd rather skip the step quietly instead of
failing the chain.
Adding your own operator
from celeryflow.conditions import register_operator
register_operator("$startswith",
lambda actual, expected: isinstance(actual, str)
and actual.startswith(expected))
Logging
import logging
from celeryflow import EventTask
EventTask.bind_logger(logging.getLogger("orders"))
Each task call will log Start, End, Retry, or Failed. The log
records carry task_id, process_id, parent_id, task_name, and a
correlation_id if the payload has one.
Development
git clone https://github.com/ChenYuTingJerry/CeleryFlow.git
cd CeleryFlow
uv sync --all-extras
# Unit tests (no broker needed):
uv run pytest -m "not integration"
# Integration tests (need Redis):
docker compose up -d
uv run pytest -m integration
docker compose down
# Lint:
uv run ruff check src tests
Project layout
celeryflow/
├── src/celeryflow/
│ ├── __init__.py public API
│ ├── app.py CeleryFlow Celery subclass
│ ├── builder.py FlowBuilder — parses and registers flows
│ ├── tasks.py FlowTask / EventTask base classes
│ ├── conditions.py condition operators
│ ├── encoding.py JSON encoder (Decimal / Enum / datetime)
│ ├── utils.py StringConvert / DictConvert
│ └── exceptions.py
├── tests/
├── docker-compose.yml Redis for integration tests
└── pyproject.toml
If you already use Celery
CeleryFlow doesn't change anything you already have. Your @app.task
definitions, broker, workers, and monitoring all keep working. The
three things CeleryFlow adds are:
CeleryFlow(...)— aCelerysubclass. Use it instead ofCelery. Takes aworker_nameand otherwise behaves the same.EventTask— aTasksubclass. Passbase=EventTaskto tasks you want in a flow. (You can also use it on every task if you just want the logging.)FlowBuilder.from_*— registers an entry task per flow. When that entry task runs, it splices the rest of the flow into the chain.
So a minimal switch looks like:
# before
from celery import Celery
app = Celery("my_worker", broker="redis://...")
@app.task
def step_a(payload): ...
@app.task
def step_b(payload): ...
# after
from celeryflow import CeleryFlow, EventTask, FlowBuilder
app = CeleryFlow("my_worker", broker="redis://...")
@app.task(base=EventTask)
def step_a(payload): ...
@app.task(base=EventTask)
def step_b(payload): ...
FlowBuilder.from_dict(app, {
"main-flows": [
{"name": "MyFlow",
"flows": [{"task": "step_a"}, {"task": "step_b"}]},
],
})
app.execute_flow("MyFlow", payload)
You can adopt it bit by bit. Flows live alongside any chain() /
group() / chord() you already build by hand.
License
MIT. See LICENSE.
Project details
Release history Release notifications | RSS feed
Download files
Download the file for your platform. If you're not sure which to choose, learn more about installing packages.
Source Distribution
Built Distribution
Filter files by name, interpreter, ABI, and platform.
If you're not sure about the file name format, learn more about wheel file names.
Copy a direct link to the current filters
File details
Details for the file celeryflow-0.1.0.tar.gz.
File metadata
- Download URL: celeryflow-0.1.0.tar.gz
- Upload date:
- Size: 17.9 kB
- Tags: Source
- Uploaded using Trusted Publishing? Yes
- Uploaded via: twine/6.1.0 CPython/3.13.12
File hashes
| Algorithm | Hash digest | |
|---|---|---|
| SHA256 |
df2b4bee1fc7d66591733c879b271949b00ad4c8771a53115eb6a92c42ce9a19
|
|
| MD5 |
650189d81a86e249ca3a72c253176830
|
|
| BLAKE2b-256 |
65e0069bc714f7b9c5e18a31b41a9737feff0cff77df0412faf6dfd204855950
|
Provenance
The following attestation bundles were made for celeryflow-0.1.0.tar.gz:
Publisher:
release.yml on ChenYuTingJerry/CeleryFlow
-
Statement:
-
Statement type:
https://in-toto.io/Statement/v1 -
Predicate type:
https://docs.pypi.org/attestations/publish/v1 -
Subject name:
celeryflow-0.1.0.tar.gz -
Subject digest:
df2b4bee1fc7d66591733c879b271949b00ad4c8771a53115eb6a92c42ce9a19 - Sigstore transparency entry: 1553794401
- Sigstore integration time:
-
Permalink:
ChenYuTingJerry/CeleryFlow@f4577f3152e22653c31200cd823bbe261fbe2a16 -
Branch / Tag:
refs/tags/v0.1.0 - Owner: https://github.com/ChenYuTingJerry
-
Access:
public
-
Token Issuer:
https://token.actions.githubusercontent.com -
Runner Environment:
github-hosted -
Publication workflow:
release.yml@f4577f3152e22653c31200cd823bbe261fbe2a16 -
Trigger Event:
push
-
Statement type:
File details
Details for the file celeryflow-0.1.0-py3-none-any.whl.
File metadata
- Download URL: celeryflow-0.1.0-py3-none-any.whl
- Upload date:
- Size: 15.5 kB
- Tags: Python 3
- Uploaded using Trusted Publishing? Yes
- Uploaded via: twine/6.1.0 CPython/3.13.12
File hashes
| Algorithm | Hash digest | |
|---|---|---|
| SHA256 |
62e1c98f4eaa8e8b89d1db009e5fa1796d1423c34a7c1797fe93b716f98d47a6
|
|
| MD5 |
8d3d7bb6303071273fefe5d7b6f920e8
|
|
| BLAKE2b-256 |
647c6158bd9f2007bc3c1bb521c2bf6119070c361dcb3178171b4b4bc743af85
|
Provenance
The following attestation bundles were made for celeryflow-0.1.0-py3-none-any.whl:
Publisher:
release.yml on ChenYuTingJerry/CeleryFlow
-
Statement:
-
Statement type:
https://in-toto.io/Statement/v1 -
Predicate type:
https://docs.pypi.org/attestations/publish/v1 -
Subject name:
celeryflow-0.1.0-py3-none-any.whl -
Subject digest:
62e1c98f4eaa8e8b89d1db009e5fa1796d1423c34a7c1797fe93b716f98d47a6 - Sigstore transparency entry: 1553794408
- Sigstore integration time:
-
Permalink:
ChenYuTingJerry/CeleryFlow@f4577f3152e22653c31200cd823bbe261fbe2a16 -
Branch / Tag:
refs/tags/v0.1.0 - Owner: https://github.com/ChenYuTingJerry
-
Access:
public
-
Token Issuer:
https://token.actions.githubusercontent.com -
Runner Environment:
github-hosted -
Publication workflow:
release.yml@f4577f3152e22653c31200cd823bbe261fbe2a16 -
Trigger Event:
push
-
Statement type: