Checklist/Docs/Implement with Celery: API RL → enqueue → RabbitMQ → job RL + lock → side effect

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

Celery pipeline: API RL → enqueue → RabbitMQ → job RL + lock → side effect
ClientFastAPIAPI rate limitEnqueuejobs row + 202RabbitMQjobs.ioWorkerjob RL + lockSide effectvendor / DBRedis RLPostgresstatus → jobs table
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
StageWhereStore
API rate limitFastAPI dependency / middlewareRedis
EnqueueFastAPI routePostgres jobs + RabbitMQ via Celery
Job rate limitCelery task (before work)Redis (shared across workers)
Entity lockCelery taskRedis SET NX EX + token release
Side effectCelery taskVendor / DB
Product statusAlways`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 = 120

Broker = 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?

  1. Rate limit first : avoid holding a lock while waiting on a vendor quota.
  2. Lock second : exclusive critical section only while doing real work.
  3. Always release in finally with token check.

Optional: take a global vendor token and a per-tenant token (two take_token calls).

---

6. Celery rate_limit vs Redis bucket

MechanismScopeUse
@task(rate_limit="30/m")Per-worker process (approximate)Extra safety net
Redis token bucketAll workersReal 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:pass

Compose 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 / .delay after 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

External