Building an Asynchronous Task Pipeline with FastAPI, Redis, and Celery: Complete Production Blueprint

Dileep Solanki

1. The Problem: The Blocking Synchronous API Trap

Modern web APIs are designed around fast, concurrent I/O operations. When built with high-performance asynchronous frameworks like FastAPI (powered by Starlette and Uvicorn), an API can comfortably serve tens of thousands of requests per second—provided those requests are lightweight database queries or fast cache lookups.

However, when a web endpoint is tasked with heavy operations—such as generating a complex PDF invoice, executing a machine learning model, resizing high-resolution image uploads, or syncing 5,000 CRM contacts—executing that logic synchronously inside the request handler blocks the Python event loop or exhausts Uvicorn worker processes. Under traffic spikes, incoming connections queue up, response times degrade from 40ms to 30 seconds, and the upstream load balancer begins returning HTTP 504 Gateway Timeout errors.

The standard enterprise solution is an Event-Driven Asynchronous Task Architecture. The FastAPI endpoint acts strictly as a lightweight Producer: it validates input, queues a task payload in an in-memory message broker (Redis), and immediately returns an HTTP 202 Accepted response alongside a tracking UUID. Independent background workers (Celery) pull tasks from Redis, execute them in isolated processes, and store the output in a persistent result backend.

2. Architectural Components

  • Producer (FastAPI): Handles authentication, payload validation via Pydantic V2, and task serialization. Responds to client requests in under 20 milliseconds.
  • Message Broker (Redis 7+): Fast, memory-resident data structure store operating as an in-memory FIFO queue for serialized JSON task messages.
  • Consumer (Celery Worker Pool): Multi-process worker pool executing tasks independently across one or more dedicated worker servers.
  • Result Backend (Redis / PostgreSQL): Persists task execution status (PENDING, STARTED, SUCCESS, FAILURE) and output data for client polling or webhook dispatches.

3. End-to-End Production Code Implementation

Step 1: Install Production Dependencies

pip install fastapi uvicorn celery redis pydantic flower

Step 2: Configure Celery Engine (celery_app.py)

from celery import Celery
import os

# Initialize Celery with Redis broker and result backend
REDIS_URL = os.getenv("REDIS_URL", "redis://localhost:6379/0")
RESULT_URL = os.getenv("RESULT_URL", "redis://localhost:6379/1")

celery_app = Celery("jioai_worker", broker=REDIS_URL, backend=RESULT_URL)

celery_app.conf.update(
    task_serializer="json",
    result_serializer="json",
    accept_content=["json"],
    timezone="UTC",
    enable_utc=True,
    task_track_started=True, # Allows client to see when worker has picked up task
    task_acks_late=True, # Worker acknowledges task ONLY after completion
    worker_prefetch_multiplier=1, # Prevents long tasks from hogging a single worker
    result_expires=86400 # Results expire from Redis after 24 hours
)

Step 3: Define Resilient Tasks with Exponential Backoff (tasks.py)

from celery_app import celery_app
import time
import logging

logger = logging.getLogger(__name__)

@celery_app.task(
    bind=True,
    max_retries=3,
    default_retry_delay=5,
    autoretry_for=(Exception,),
    retry_backoff=True, # Exponential backoff: 5s, 10s, 20s
    retry_backoff_max=60,
    retry_jitter=True # Random jitter to prevent thundering herd on external APIs
)
def generate_customer_report(self, customer_id: int, report_format: str):
    logger.info("Starting task for customer " + str(customer_id))
    try:
        # Simulate heavy computational task
        time.sleep(4)
        result = {
            "customer_id": customer_id,
            "format": report_format,
            "download_url": "https://cdn.jioai.tech/reports/" + str(customer_id) + "_export." + report_format,
            "records_processed": 4820
        }
        return result
    except Exception as exc:
    logger.error("Task failed: " + str(exc) + ". Retrying...")
        raise exc

Step 4: FastAPI Router with Polling Endpoints (main.py)

from fastapi import FastAPI, HTTPException, status
from pydantic import BaseModel, Field
from celery_app import celery_app
from tasks import generate_customer_report
from celery.result import AsyncResult

app = FastAPI(title="JioAI Asynchronous Task Service")

class ReportRequest(BaseModel):
    customer_id: int = Field(gt=0, description="Valid Customer ID")
    report_format: str = Field(default="pdf", pattern="^(pdf|csv|xlsx)$")

class TaskResponse(BaseModel):
    task_id: str
    status: str

@app.post("/reports/generate", status_code=status.HTTP_202_ACCEPTED, response_model=TaskResponse)
async def dispatch_report_task(payload: ReportRequest):
    # Dispatch task asynchronously to Redis
    task = generate_customer_report.delay(payload.customer_id, payload.report_format)
    return TaskResponse(task_id=task.id, status="QUEUED")

@app.get("/reports/status/{task_id}")
async def get_task_status(task_id: str):
    res = AsyncResult(task_id, app=celery_app)
    response_data = {
        "task_id": task_id,
        "state": res.state
    }
    if res.state == "SUCCESS":
        response_data["result"] = res.result
    elif res.state == "FAILURE":
        response_data["error"] = str(res.result)
    return response_data

4. Running and Scaling in Production

To run this pipeline locally or inside Docker containers, execute three independent services:

# Terminal 1: Launch Redis Message Broker
docker run -d -p 6379:6379 --name redis-broker redis:7-alpine

# Terminal 2: Start Celery Worker Pool with 4 concurrent processes
celery -A celery_app worker --loglevel=info --concurrency=4

# Terminal 3: Start FastAPI Uvicorn Web Server
uvicorn main:app --host 0.0.0.0 --port 8000 --reload

# Terminal 4: Launch Real-Time Monitoring UI (Flower)
celery -A celery_app flower --port=5555

5. Persistent Database Integration with SQLAlchemy and Dead-Letter Queues

In enterprise architectures, tasks cannot exist in isolation from relational application state. When an asynchronous worker completes a report or processes a batch of invoices, the status and metadata must be atomically written to a relational database like PostgreSQL using SQLAlchemy 2.0.

Furthermore, production queues must account for "Poison Pill" tasks: malformed requests that trigger unhandled runtime exceptions. If a task crashes a worker process, Celery must route the failed payload to an isolated Dead-Letter Queue (DLQ) rather than allowing it to endlessly retry and block worker capacity.

# database.py - Asynchronous SQLAlchemy Session Manager
from sqlalchemy.ext.asyncio import create_async_engine, async_sessionmaker, AsyncSession
from sqlalchemy.orm import DeclarativeBase, Mapped, mapped_column
from sqlalchemy import String, Integer, DateTime, func

DATABASE_URL = "postgresql+asyncpg://app_user:secret@localhost:5432/app_db"
engine = create_async_engine(DATABASE_URL, pool_size=20, max_overflow=10)
AsyncSessionLocal = async_sessionmaker(engine, expire_on_commit=False)

class Base(DeclarativeBase):
    pass

class TaskAuditLog(Base):
    __tablename__ = "task_audit_logs"
    id: Mapped[int] = mapped_column(primary_key=True)
    task_uuid: Mapped[str] = mapped_column(String(64), index=True)
    customer_id: Mapped[int] = mapped_column(Integer, index=True)
    status: Mapped[str] = mapped_column(String(32), default="PENDING")
    created_at: Mapped[DateTime] = mapped_column(DateTime, server_default=func.now())

Dead-Letter Queue Routing in Celery

# Configure specialized DLQ routing in celery_app.py
from kombu import Queue, Exchange

default_exchange = Exchange("default", type="direct")
dlq_exchange = Exchange("dlq", type="direct")

celery_app.conf.task_queues = (
    Queue("default", default_exchange, routing_key="default"),
    Queue("dead_letter", dlq_exchange, routing_key="dead_letter"),
)
celery_app.conf.task_default_queue = "default"
celery_app.conf.task_default_exchange = "default"
celery_app.conf.task_default_routing_key = "default"

5. Production Hardening Checklist

Reliability Rules:

  • Always set task_acks_late=True: By default, Celery acknowledges a task the millisecond it pulls it off the queue. If the worker server crashes during processing, the task is lost forever. Setting late acknowledgments ensures tasks re-queue upon worker termination.

  • Set worker_prefetch_multiplier=1: Celery by default prefetches 4 tasks per process. If one worker receives a 20-minute video transcoding task, it hoards 3 other fast tasks behind it. Setting prefetch to 1 ensures fair work distribution.

  • Enable Redis Persistence: Configure Redis with Append-Only File (AOF) logging (appendonly yes) so queued jobs survive power failures or container restarts.


3/related/default