Guide 17
Implement with Celery: API RL → enqueue → RabbitMQ → job RL + lock → side effect
Audience: FastAPI + Celery + RabbitMQ teams implementing the standard pipeline. Related: 16 Locks and rate limits · 05 Celery example
Pipeline
text
Client
│
▼
FastAPI ── API rate limit (Redis) ──► 429 or continue
│
├── insert jobs row (pending)
└── Celery apply_async ──► RabbitMQ queue
│
▼
Celery worker
│
┌─────────────┼─────────────┐
▼ ▼ ▼
job rate limit entity lock load job
(Redis) (Redis) (DB)
│ │ │
└─────────────┴──────► side effect
│
mark succeeded / failed| Stage | Where | Store |
|---|---|---|
| API rate limit | FastAPI dependency / middleware | Redis |
| Enqueue | FastAPI route | Postgres jobs + RabbitMQ via Celery |
| Job rate limit | Celery task (before work) | Redis (shared across workers) |
| Entity lock | Celery task | Redis SET NX EX + token release |
| Side effect | Celery task | Vendor / DB |
| Product status | Always | `jobs` table, not Flower |
---
1. Dependencies and settings
python
# PSEUDOCODE : requirements
# celery[redis] # redis extra only if you use Redis result backend; broker is RabbitMQ
# redis
# fastapi, pydantic-settings, sqlalchemy/asyncpg, ...
# PSEUDOCODE : settings
class Settings(BaseSettings):
database_url: str
celery_broker_url: str = "amqps://user:pass@rabbitmq:5671//"
redis_url: str = "redis://redis:6379/0" # locks + rate limits only
api_enqueue_limit: int = 30
api_enqueue_window_sec: int = 60
vendor_rate_per_sec: float = 5.0
vendor_burst: int = 10
lock_ttl_sec: int = 120Broker = RabbitMQ. Redis is not the job broker.
---
2. Redis helpers (locks + rate limits)
python
# PSEUDOCODE : app/core/redis_limits.py
import time
import uuid
from redis import Redis
RELEASE_LUA = """
if redis.call("get", KEYS[1]) == ARGV[1] then
return redis.call("del", KEYS[1])
else
return 0
end
"""
# Simple token bucket (production: use a well-tested Lua bucket)
BUCKET_LUA = """
local key = KEYS[1]
local rate = tonumber(ARGV[1])
local burst = tonumber(ARGV[2])
local now = tonumber(ARGV[3])
local data = redis.call("hmget", key, "tokens", "ts")
local tokens = tonumber(data[1])
local ts = tonumber(data[2])
if tokens == nil then
tokens = burst
ts = now
end
local delta = math.max(0, now - ts)
tokens = math.min(burst, tokens + delta * rate)
local ok = 0
if tokens >= 1 then
tokens = tokens - 1
ok = 1
end
redis.call("hmset", key, "tokens", tokens, "ts", now)
redis.call("expire", key, 3600)
return ok
"""
def get_redis(url: str) -> Redis:
return Redis.from_url(url, decode_responses=True)
def api_allow(r: Redis, key: str, limit: int, window_sec: int) -> bool:
n = r.incr(key)
if n == 1:
r.expire(key, window_sec)
return n <= limit
def take_token(r: Redis, key: str, rate_per_sec: float, burst: int) -> bool:
return r.eval(BUCKET_LUA, 1, key, rate_per_sec, burst, time.time()) == 1
def acquire_lock(r: Redis, key: str, token: str, ttl_sec: int) -> bool:
return bool(r.set(key, token, nx=True, ex=ttl_sec))
def release_lock(r: Redis, key: str, token: str) -> None:
r.eval(RELEASE_LUA, 1, key, token)Use sync Redis in Celery prefork workers; use async Redis only if your worker model is async (Taskiq) or you wrap calls carefully.
---
3. Celery app (RabbitMQ)
python
# PSEUDOCODE : app/workers/celery_app.py
from celery import Celery
from app.core.settings import settings
celery_app = Celery("app", broker=settings.celery_broker_url)
celery_app.conf.update(
task_serializer="json",
accept_content=["json"],
result_serializer="json",
task_acks_late=True,
worker_prefetch_multiplier=1,
task_default_queue="jobs.default",
task_routes={
"app.workers.tasks.sync.*": {"queue": "jobs.io"},
},
broker_connection_retry_on_startup=True,
)
celery_app.autodiscover_tasks(["app.workers.tasks"])---
4. FastAPI: API rate limit → enqueue
python
# PSEUDOCODE : app/api/deps.py
from fastapi import Depends, HTTPException, Request
from app.core.redis_limits import get_redis, api_allow
from app.core.settings import settings
def get_r():
return get_redis(settings.redis_url)
def rate_limit_enqueue(
request: Request,
user=Depends(get_current_user),
r=Depends(get_r),
):
# per-user (and optionally also per-IP)
key = f"rl:api:enqueue:user:{user.id}"
if not api_allow(r, key, settings.api_enqueue_limit, settings.api_enqueue_window_sec):
raise HTTPException(
status_code=429,
detail="Too many job submissions",
headers={"Retry-After": str(settings.api_enqueue_window_sec)},
)python
# PSEUDOCODE : app/api/routes/sync.py
from fastapi import APIRouter, Depends
from app.workers.tasks.sync import sync_account
from app.api.deps import rate_limit_enqueue, get_db
router = APIRouter(prefix="/sync", tags=["sync"])
@router.post("/accounts/{account_id}", status_code=202, dependencies=[Depends(rate_limit_enqueue)])
def start_sync(account_id: str, db=Depends(get_db), user=Depends(get_current_user)):
# optional: unique key so double-click does not double-publish
idem = f"sync-account:{account_id}:{user.id}"
existing = db.find_job_by_idempotency(idem)
if existing:
return {"job_id": existing.id, "status": existing.status}
job = db.insert_job(
type="sync_account",
status="pending",
entity_id=account_id,
idempotency_key=idem,
)
# Celery → RabbitMQ
sync_account.apply_async(
args=[str(job.id), account_id],
queue="jobs.io",
)
return {"job_id": job.id, "status": "pending"}Prefer outbox if you need commit-safe publish (04); apply_async right after insert is the simple path.
---
5. Celery task: job rate limit + entity lock → side effect
python
# PSEUDOCODE : app/workers/tasks/sync.py
import uuid
from celery.exceptions import MaxRetriesExceededError
from app.workers.celery_app import celery_app
from app.core.settings import settings
from app.core.redis_limits import (
get_redis,
take_token,
acquire_lock,
release_lock,
)
from app.db import session_scope
from app.services.vendor import call_vendor_sync
@celery_app.task(
bind=True,
name="app.workers.tasks.sync.sync_account",
max_retries=25,
acks_late=True,
soft_time_limit=90,
time_limit=120,
)
def sync_account(self, job_id: str, account_id: str) -> None:
r = get_redis(settings.redis_url)
# --- job / vendor rate limit (shared across all workers) ---
rl_key = f"rl:vendor:sync:{account_id}" # or global: rl:vendor:sync
if not take_token(
r,
rl_key,
rate_per_sec=settings.vendor_rate_per_sec,
burst=settings.vendor_burst,
):
# re-queue with delay; does not burn forever if max_retries set
raise self.retry(countdown=2 + (self.request.retries % 5))
# --- entity lock (one sync per account at a time) ---
token = str(uuid.uuid4())
lock_key = f"lock:account:{account_id}"
if not acquire_lock(r, lock_key, token, ttl_sec=settings.lock_ttl_sec):
raise self.retry(countdown=10 + (self.request.retries % 10))
try:
with session_scope() as db:
job = db.get_job(job_id)
if job is None:
return
if job.status == "succeeded":
return # idempotent
db.mark_running(job_id)
try:
call_vendor_sync(account_id) # side effect
db.mark_succeeded(job_id)
except PermanentVendorError as e:
db.mark_failed(job_id, str(e))
# do not retry permanent errors
return
except TransientVendorError as e:
db.bump_attempt(job_id, str(e))
raise self.retry(
exc=e,
countdown=min(600, 2 ** self.request.retries),
)
finally:
release_lock(r, lock_key, token)Why this order?
- Rate limit first : avoid holding a lock while waiting on a vendor quota.
- Lock second : exclusive critical section only while doing real work.
- Always release in
finallywith token check.
Optional: take a global vendor token and a per-tenant token (two take_token calls).
---
6. Celery rate_limit vs Redis bucket
| Mechanism | Scope | Use |
|---|---|---|
@task(rate_limit="30/m") | Per-worker process (approximate) | Extra safety net |
| Redis token bucket | All workers | Real multi-worker fairness |
python
# Optional second line of defense (not enough alone with many workers)
@celery_app.task(bind=True, rate_limit="30/m")
def sync_account(self, job_id: str, account_id: str):
...Prefer Redis for production multi-replica workers.
---
7. Run processes
bash
# API
uvicorn app.main:app --host 0.0.0.0 --port 8000
# Workers consuming the IO queue
celery -A app.workers.celery_app.celery_app worker -Q jobs.io,jobs.default -c 4
# Optional Flower (private + auth) : not product status
celery -A app.workers.celery_app.celery_app flower --basic_auth=user:passCompose services: api, worker, rabbitmq, redis, db.
---
8. Testing sketch
python
# PSEUDOCODE
def test_enqueue_rate_limited(client, redis):
for _ in range(30):
assert client.post("/sync/accounts/a1").status_code == 202
assert client.post("/sync/accounts/a1").status_code == 429
def test_task_retries_when_lock_held(redis):
acquire_lock(redis, "lock:account:a1", "other", 60)
with pytest.raises(Retry):
sync_account.push_request(retries=0).run("job-1", "a1") # or mock retry
def test_success_releases_lock(redis, db):
sync_account.run(str(job.id), "a1")
assert redis.get("lock:account:a1") is None
assert db.get_job(job.id).status == "succeeded"---
9. Production checklist (Celery-specific)
- [ ] Broker URL is RabbitMQ (
amqps://in prod) - [ ] Redis used only for RL + locks (and optional result backend)
- [ ] API enqueue dependency returns 429 when over limit
- [ ]
apply_async/.delayafter job row insert (or outbox) - [ ] Task uses acks_late, JSON, routed queue
- [ ] Redis job rate limit before side effect
- [ ] Entity lock with token-safe release in
finally - [ ] Lock busy / RL miss →
self.retry(countdown=…)with max_retries - [ ] Permanent errors →
mark_failed, no retry - [ ] Status for users from `GET /jobs/{id}`
- [ ] Metrics:
api_429,lock_busy_retry,rl_requeue, runtime
---
10. Minimal file layout
text
app/
core/settings.py
core/redis_limits.py
api/deps.py
api/routes/sync.py
workers/celery_app.py
workers/tasks/sync.py
db/...---
See also
- 16 Locks and rate limits : concepts
- 05 Celery + RabbitMQ : app wiring
- 09 Errors / DLQ
- 15 Observability & Flower
- 03 Job lifecycle