If your API ever sends an email, generates a PDF, calls a slow third party service, or processes a file upload, you have almost certainly made your users wait for something they did not need to watch happen live. That wait is not a performance problem you fix with faster code. It is an architecture problem, and the fix is a background task queue.
This guide builds that queue completely with FastAPI, Celery, and Redis, running in Docker. Every file in the project gets shown, in the order you would actually create them, with an explanation of what it does before the code and a breakdown of anything non obvious after it.
Why Your API Should Never Do Slow Work In The Request Cycle
A request handler that blocks on a slow operation ties up a worker process for the entire duration of that operation. Do this often enough under load and your API becomes unresponsive even for requests that have nothing to do with the slow task. The fix is to separate accepting the work from doing the work. Your API accepts the request, records what needs to happen, hands it off to something else, and responds immediately. That something else is Celery.
The Three Moving Parts
FastAPI is your web layer. It accepts HTTP requests, validates input, and returns responses quickly, and should never block on slow work.
Celery is your task execution engine. It runs Python functions outside the request response cycle, on separate worker processes.
Redis plays two roles at once. As the broker, it is the queue holding tasks waiting to be picked up by a worker. As the result backend, it stores a task's outcome so something else can check on it later.
Docker ties these together so your API, your worker, your scheduler, and Redis each run as separate, reproducible services that start with one command.
Project Structure
app/
core/
config.py
database.py
celery_app.py
jobs/
models.py
schemas.py
repository.py
service.py
tasks.py
routes.py
main.py
Dockerfile
docker-compose.yml
requirements.txt
.env
The repository only queries and persists data, without committing. A narrow service applies business rules, also without committing. One orchestrating service per use case owns the transaction, commits it, and is the only place that triggers a Celery task. Routes parse input, call a single service method, and translate exceptions into HTTP responses.
requirements.txt
This file pins every Python package the project depends on, so the Docker image you build today installs the exact same versions a year from now.
fastapi==0.115.0
uvicorn[standard]==0.30.6
celery[redis]==5.4.0
redis==5.0.8
sqlalchemy==2.0.35
psycopg2-binary==2.9.9
pydantic==2.9.2
pydantic-settings==2.5.2
python-dotenv==1.0.1
alembic==1.13.3
fastapi and uvicorn run the web layer, uvicorn being the actual server process that executes your FastAPI app. celery[redis] installs Celery along with the extra dependencies it needs specifically to talk to Redis as a broker, without the [redis] extra you would need to install redis separately anyway, which is listed here too since your own code also talks to Redis directly for things like idempotency locks. sqlalchemy and psycopg2-binary give you the ORM and the actual PostgreSQL driver it uses underneath. pydantic-settings is what makes the configuration file below possible, reading environment variables into a typed settings object. python-dotenv lets that settings object load values from a local .env file during development. alembic handles database migrations, not shown in this guide but assumed present in any real project with a Job table.
.env
This file holds configuration that changes between environments and should never be committed to version control.
DATABASE_URL=postgresql+psycopg2://app:app@db:5432/app
REDIS_URL=redis://redis:6379/0
ENVIRONMENT=production
SECRET_KEY=change-this-to-a-random-value
DATABASE_URL points at the db service by its Docker Compose service name rather than localhost, since containers on the same Compose network reach each other by service name, not by the host machine's loopback address. REDIS_URL works the same way, pointing at the redis service and selecting database index 0 inside it. ENVIRONMENT is a simple flag your own code can branch on if it needs to behave differently in development versus production. SECRET_KEY is a placeholder for whatever secret your app needs for things like signing tokens, and should be a long random value in any real deployment, never the literal text shown here.
app/core/config.py
This module turns environment variables into a single typed object the rest of the app imports, instead of scattering os.environ.get() calls throughout the codebase.
# app/core/config.py
from pydantic_settings import BaseSettings
class Settings(BaseSettings):
DATABASE_URL: str
REDIS_URL: str
ENVIRONMENT: str = "development"
SECRET_KEY: str
class Config:
env_file = ".env"
settings = Settings()
BaseSettings from Pydantic reads matching environment variables automatically, falling back to the .env file named in Config.env_file when a variable is not already set in the environment. Each field is typed, so if DATABASE_URL is missing entirely, your app fails immediately at startup with a clear error, rather than failing later with a confusing connection error the first time something tries to use it. The settings instance created at the bottom is imported everywhere else that needs configuration, including the database setup and the Celery app itself.
app/core/database.py
This module creates the actual connection to PostgreSQL and the machinery FastAPI and Celery both use to get a database session.
# app/core/database.py
from sqlalchemy import create_engine
from sqlalchemy.orm import sessionmaker, declarative_base
from app.core.config import settings
engine = create_engine(settings.DATABASE_URL, pool_pre_ping=True)
SessionLocal = sessionmaker(autocommit=False, autoflush=False, bind=engine)
Base = declarative_base()
def get_db():
db = SessionLocal()
try:
yield db
finally:
db.close()
engine is the actual connection pool to the database, created once when this module first loads. pool_pre_ping=True makes SQLAlchemy test a connection before handing it out, which matters a great deal in Docker, where the database container can restart or briefly drop connections without your application process restarting alongside it. SessionLocal is a factory, calling it produces a new session, it is not itself a session. Base is the class every model inherits from, which is how SQLAlchemy knows which Python classes map to which database tables. get_db is a generator function FastAPI calls through Depends, it opens a session, hands it to the route for the duration of that one request, and guarantees the session is closed afterward regardless of whether the request succeeded or raised an exception. Celery tasks cannot use get_db this way since there is no request lifecycle to hook into, which is why tasks call SessionLocal() directly instead, as you will see in tasks.py.
app/core/celery_app.py
This module creates the Celery application itself and configures how it behaves, and it is the single file both the FastAPI process and the worker process import to agree on the same broker, backend, and settings.
# app/core/celery_app.py
from celery import Celery
from celery.schedules import crontab
from app.core.config import settings
celery_app = Celery(
"worker",
broker=settings.REDIS_URL,
backend=settings.REDIS_URL,
include=["app.jobs.tasks"],
)
celery_app.conf.update(
task_serializer="json",
result_serializer="json",
accept_content=["json"],
timezone="UTC",
enable_utc=True,
task_acks_late=True,
worker_prefetch_multiplier=1,
task_reject_on_worker_lost=True,
result_expires=3600,
task_routes={
"app.jobs.tasks.generate_report": {"queue": "reports"},
"app.jobs.tasks.send_notification_email": {"queue": "emails"},
},
)
celery_app.conf.beat_schedule = {
"cleanup-expired-jobs-daily": {
"task": "app.jobs.tasks.cleanup_expired_jobs",
"schedule": crontab(hour=2, minute=0),
},
}
Every parameter in the Celery constructor
"worker" is the first positional argument, and it is the name of the Celery application itself, not a reference to the worker process you run later. It is used internally to namespace task names and shows up as an identifier in logs and in Flower's dashboard. You could name it "myapp" and nothing else in this guide would change, the celery -A ... worker command still works the same way.
broker=settings.REDIS_URL is the message queue a task actually lands in the moment you call .delay() or .apply_async(). Setting this to your Redis URL means Celery pushes a serialized message describing the task, its name, its arguments, and some metadata, onto a Redis list, and a running worker process is what watches that list and pulls messages off it. Without a reachable broker, calling .delay() from your FastAPI service either fails outright with a connection error or queues a message nothing is listening for.
backend=settings.REDIS_URL is a separate concern from the broker, even though this example points both at the same Redis instance. The backend is where Celery stores what happened after a task finished, its return value, its final state, and a traceback if it failed. You only need this configured if something later checks on a task using its id through AsyncResult. If nothing in your system ever asks whether a task finished, which is common for fire and forget work like sending an email, you can set ignore_result=True on that specific task and skip the backend writes entirely for it.
include=["app.jobs.tasks"] tells Celery which modules to import the moment the application starts, specifically so the @celery_app.task decorated functions inside them get registered with this Celery app. This line is easy to misunderstand. celery_app.py itself never imports tasks.py anywhere in its own code, so when you start a worker with celery -A app.core.celery_app.celery_app worker, Celery would otherwise only know about this one file, not any task you have written elsewhere. The include list is what makes Celery go import those modules on startup and discover every decorated task inside them, registering their names so the worker can execute them when a message arrives asking for app.jobs.tasks.generate_report by that exact string. Forget to add a new tasks module here later and the worker logs a confusing "received unregistered task" error even though the function clearly exists in your code.
Every setting in conf.update
task_serializer and result_serializer set to "json" instead of Celery's older pickle default avoid deserializing arbitrary Python objects on the worker side, which is both a security risk and a source of subtle bugs whenever your models change shape between when a task was queued and when it runs. accept_content=["json"] locks the worker to only accept json formatted messages, refusing anything else, which closes the same security gap from the receiving end.
timezone="UTC" and enable_utc=True make every scheduled time, including the beat schedule below, unambiguous regardless of what timezone the server hosting the worker happens to be in.
task_acks_late=True changes when a task is considered handled. By default Celery acknowledges a task, meaning it removes it from the queue, the moment a worker picks it up, before the task has actually run. With this set to true, acknowledgment only happens after the task finishes successfully, so if the worker process crashes mid task, the message is never acknowledged and gets redelivered to another worker instead of silently disappearing.
worker_prefetch_multiplier=1 stops a single worker from grabbing a large batch of queued tasks upfront. The default prefetch behavior can cause one worker to hoard a dozen tasks while sitting idle workers starve, which matters once your tasks vary in how long they take to run.
task_reject_on_worker_lost=True works alongside acks_late, explicitly telling Celery to requeue a task if the worker executing it is terminated or lost partway through, rather than leaving that task's fate ambiguous.
result_expires=3600 automatically deletes stored task results from Redis after one hour, so results you checked once and no longer need do not accumulate in memory indefinitely.
task_routes sends different kinds of work to different named queues. This is the detail that saves you the day a two minute report generation task stops a time sensitive notification email from being sent for two minutes, because they are no longer competing for the same worker's attention.
Where the beat schedule code belongs
The celery_app.conf.beat_schedule block shown above belongs in this same file, app/core/celery_app.py, appended directly after the conf.update() call. Celery Beat reads its schedule from the exact Celery app instance you point it at when you run celery -A app.core.celery_app.celery_app beat, so if the schedule lived in some other file that celery_app.py never imports, Beat would start up with an empty schedule and the cleanup task would simply never fire, with no error telling you why. If your list of scheduled tasks grows large enough to feel cluttered here, you can move the dictionary into its own app/core/beat_schedule.py and import it at the bottom of celery_app.py, but it must still end up assigned to celery_app.conf.beat_schedule one way or another.
app/jobs/models.py
This defines the actual database table backing a job, using SQLAlchemy's declarative style.
# app/jobs/models.py
from sqlalchemy import Column, Integer, String, DateTime, func
from app.core.database import Base
class Job(Base):
__tablename__ = "jobs"
id = Column(Integer, primary_key=True, index=True)
status = Column(String, default="pending", nullable=False)
result_path = Column(String, nullable=True)
created_at = Column(DateTime(timezone=True), server_default=func.now())
updated_at = Column(DateTime(timezone=True), onupdate=func.now())
status tracks the job through its lifecycle, pending, processing, completed, or failed, and is the field the idempotency guard in the task checks. result_path stays null until the task finishes, then holds wherever the generated report actually ended up. created_at is set by the database itself at insert time through server_default, and updated_at refreshes automatically on every update through onupdate, so neither needs to be set manually anywhere in the application code.
app/jobs/schemas.py
This defines the Pydantic models that validate incoming requests and shape outgoing responses, kept separate from the database model above so your API's public shape can evolve independently of your table structure.
# app/jobs/schemas.py
from datetime import datetime
from pydantic import BaseModel
class JobCreate(BaseModel):
report_type: str
class JobRead(BaseModel):
id: int
status: str
result_path: str | None
created_at: datetime
class Config:
from_attributes = True
JobCreate defines exactly what a client must send to create a job, here just the kind of report being requested, and FastAPI validates incoming JSON against it automatically before your route code ever runs. JobRead defines what gets sent back, deliberately excluding internal fields you might not want exposed. from_attributes = True is what lets you return a SQLAlchemy Job instance directly from a route and have Pydantic read its attributes to build the response, rather than requiring you to manually convert it to a dictionary first.
app/jobs/repository.py
This is the only place in the application that writes raw SQLAlchemy queries for jobs, and it never calls commit(), that responsibility belongs entirely to the service layer above it.
# app/jobs/repository.py
from sqlalchemy.orm import Session
from app.jobs.models import Job
from app.jobs.schemas import JobCreate
class JobRepository:
def __init__(self, db: Session):
self.db = db
def create(self, payload: JobCreate) -> Job:
job = Job(status="pending", **payload.model_dump())
self.db.add(job)
self.db.flush()
return job
def get(self, job_id: int) -> Job | None:
return self.db.get(Job, job_id)
def update_status(self, job_id: int, status: str, result_path: str | None = None) -> None:
job = self.db.get(Job, job_id)
job.status = status
if result_path is not None:
job.result_path = result_path
create builds a new Job from the validated input and calls flush() rather than commit(). Flushing sends the insert statement to the database and populates the object's generated id, without ending the transaction, which lets the calling service decide when the transaction actually finishes. get is a plain lookup by primary key. update_status fetches the row and mutates it in place, relying on SQLAlchemy's unit of work to know the object is dirty and needs writing once something else commits.
app/jobs/tasks.py
This is where the actual background work lives, and it is intentionally the only file that talks to both Celery and the database directly inside a task function.
# app/jobs/tasks.py
from celery.exceptions import SoftTimeLimitExceeded
from app.core.celery_app import celery_app
from app.core.database import SessionLocal
from app.jobs.repository import JobRepository
@celery_app.task(
bind=True,
max_retries=5,
soft_time_limit=50,
time_limit=60,
)
def generate_report(self, job_id: int):
db = SessionLocal()
repo = JobRepository(db)
try:
job = repo.get(job_id)
if job.status == "completed":
return
repo.update_status(job_id, "processing")
db.commit()
result_path = build_report_file(job_id)
repo.update_status(job_id, "completed", result_path=result_path)
db.commit()
except SoftTimeLimitExceeded:
repo.update_status(job_id, "failed")
db.commit()
except Exception as exc:
db.rollback()
countdown = 2 ** self.request.retries * 10
raise self.retry(exc=exc, countdown=countdown)
finally:
db.close()
Walking through generate_report line by line
bind=True in the decorator is why self appears as the function's first parameter. This self is not an instance of a class you wrote, it is the Celery task instance itself, and binding it gives the function access to things like self.request, which carries metadata about the current execution including how many times it has already been retried, and self.retry(), used further down. job_id: int is deliberately the only piece of data passed in, a plain integer rather than the Job object itself, because Celery serializes task arguments, and an ORM object serialized at the moment a task was queued could be stale by the time it actually runs, possibly minutes later if the queue is backed up.
db = SessionLocal() opens a brand new, independent database session scoped to this one execution of the task. Celery tasks run in separate worker processes from the FastAPI process that queued them, so there is no session to inherit, this line manufactures one from scratch every time the task runs. repo = JobRepository(db) wraps that session in the same repository class used elsewhere, so the actual SQL logic for fetching and updating a job stays in exactly one place in the codebase.
Inside the try block, job = repo.get(job_id) fetches a fresh row from the database rather than trusting any data that might have been passed in, which is the fresh fetch pattern mentioned earlier, guaranteeing the task works with current state. if job.status == "completed": return is the idempotency guard, the single most important line in this function. If this exact task body somehow runs twice, because Celery redelivered an unacknowledged message after a worker crash, or a retry fired after the work had in fact already finished, this check stops it from redoing the work and potentially overwriting a valid result. In production you would also want to guard against job being None here, in case the row was deleted between the task being queued and actually running, which this simplified version leaves out to keep the flow readable.
repo.update_status(job_id, "processing") followed immediately by db.commit() marks the job as actively in progress and commits right away, so if a client checks the job's status through an API endpoint while this task is still running, they see "processing" rather than a stale "pending." result_path = build_report_file(job_id) stands in for whatever the actual slow work is, generating a PDF, aggregating data, calling an external service. It is deliberately pulled out as its own function rather than inlined here, because this task's job is orchestration, fetching, marking status, delegating to the real logic, and marking status again, not the report building logic itself. The final repo.update_status(...) and db.commit() pair records the finished state.
except SoftTimeLimitExceeded catches a specific exception Celery raises inside the task itself once it has run past the soft_time_limit set in the decorator, 50 seconds here. This is Celery's way of giving a slow task one last chance to clean up, in this case marking the job failed and committing that, before the harder time_limit of 60 seconds kills the process outright a few seconds later with no further opportunity to react. Notice this branch does not retry, a task that has already blown through its time budget likely has a structural problem that simply running it again will not fix, so it is marked failed outright instead of requeued.
except Exception as exc catches everything else, a network timeout calling some external service, an unexpected bug, a database hiccup. db.rollback() undoes any uncommitted changes from this particular attempt, so a half finished transaction never lingers. countdown = 2 ** self.request.retries * 10 computes exponential backoff using how many times this task has already been retried, 10 seconds after the first failure, 20 after the second, 40 after the third, and so on, so a repeatedly failing dependency does not get hammered at a fixed interval. raise self.retry(exc=exc, countdown=countdown) hands control back to Celery, which reschedules the task to run again after that delay, up to the max_retries=5 ceiling set in the decorator. Raising it, rather than simply calling self.retry(...) on its own line, is deliberate, it halts execution of the current attempt immediately instead of letting the function fall through past the except block.
finally: db.close() runs no matter how the function exits, success, a soft time limit, an exception, or a retry, and releases the database connection back to the pool. Skipping this in a worker that processes many tasks in sequence slowly leaks open connections until the database eventually refuses new ones.
app/jobs/service.py
This is the orchestrating service for the one use case of creating a report job, and it is the only place in the codebase allowed to both commit the transaction and dispatch the Celery task.
# app/jobs/service.py
from sqlalchemy.orm import Session
from app.jobs.repository import JobRepository
from app.jobs.schemas import JobCreate
from app.jobs.tasks import generate_report
class CreateReportJobService:
def __init__(self, db: Session):
self.db = db
self.repo = JobRepository(db)
def execute(self, payload: JobCreate):
job = self.repo.create(payload)
self.db.commit()
self.db.refresh(job)
generate_report.delay(job.id)
return job
The ordering here matters more than it looks. db.commit() happens before generate_report.delay(job.id). If the task were dispatched first, there is a real chance a worker picks it up and queries the database for a row that is not there yet, because the transaction holding it has not actually been written to disk. Committing first, then dispatching, guarantees the worker always finds the row it expects.
app/jobs/routes.py
This exposes the use case over HTTP, and does essentially nothing besides parsing input and handing off to the service.
# app/jobs/routes.py
from fastapi import APIRouter, Depends
from sqlalchemy.orm import Session
from app.core.database import get_db
from app.jobs.schemas import JobCreate, JobRead
from app.jobs.service import CreateReportJobService
router = APIRouter(prefix="/jobs", tags=["jobs"])
@router.post("", response_model=JobRead, status_code=201)
def create_job(payload: JobCreate, db: Session = Depends(get_db)):
service = CreateReportJobService(db)
job = service.execute(payload)
return job
Depends(get_db) is where the session created in database.py actually enters the request, FastAPI calls that generator function, hands the route the session it yields, and closes it automatically once the response is sent. The route itself contains no business logic at all, it exists purely to translate an HTTP request into a call to CreateReportJobService.execute.
app/jobs/report_builder.py
This is where the actual report generation logic lives, kept separate from the Celery task so it can be tested or reused without touching anything task related.
import os
from datetime import datetime, timezone
from reportlab.pdfgen import canvas
REPORTS_DIR = "generated_reports"
def build_report_file(job_id: int) -> str:
os.makedirs(REPORTS_DIR, exist_ok=True)
file_path = os.path.join(REPORTS_DIR, f"report_{job_id}.pdf")
timestamp = datetime.now(timezone.utc).isoformat()
pdf = canvas.Canvas(file_path)
pdf.drawString(100, 750, f"Report for job {job_id}")
pdf.drawString(100, 730, f"Generated at {timestamp}")
pdf.save()
return file_pathos.makedirs(REPORTS_DIR, exist_ok=True) creates the output folder the first time this runs and does nothing on every call after that, exist_ok=True is what stops it from raising an error if the folder is already there. file_path builds a unique name per job using the job's id, so two different jobs never overwrite each other's output. canvas.Canvas(file_path) from the reportlab library opens a new PDF document at that path, drawString places text at an x, y coordinate on the page, and pdf.save() writes the finished file to disk. The function returns the path it just wrote, which is exactly what tasks.py stores on the job record afterward.
In a real system this function would pull actual data relevant to the report, querying the database for whatever the report is summarizing, rather than writing two lines of placeholder text, but the shape stays the same regardless of what the report actually contains.
app/main.py
This is the entry point that assembles the FastAPI application and wires in the routes defined elsewhere.
# app/main.py
from fastapi import FastAPI
from app.jobs.routes import router as jobs_router
app = FastAPI(title="Report Service")
app.include_router(jobs_router)
app is the actual ASGI application object uvicorn runs. include_router mounts every route defined in jobs/routes.py onto the main app, under whatever prefix that router declared, /jobs in this case. A real project would include additional routers here as it grows, each still following the same pattern of staying thin and delegating to a service.
Dockerfile
This defines how to build the image that both the web service and the worker service run from, since they share the exact same codebase and dependencies, only the command run at the end differs between them.
FROM python:3.12-slim
WORKDIR /app
COPY requirements.txt .
RUN pip install --no-cache-dir -r requirements.txt
COPY . .
CMD ["uvicorn", "app.main:app", "--host", "0.0.0.0", "--port", "8000"]
FROM python:3.12-slim picks a minimal Python base image, small enough to keep build times and image size reasonable without needing every system package a full image would include. WORKDIR /app sets the working directory inside the container for every instruction that follows. COPY requirements.txt . copies only the dependency list first, before the rest of the code, which is a deliberate ordering, Docker caches each layer, so as long as requirements.txt has not changed, Docker reuses the cached layer from the previous build instead of reinstalling every package again, which is what RUN pip install does on the line after it. COPY . . then copies the rest of the application code in, happening after the dependency install specifically so that changing application code does not invalidate the dependency installation cache. CMD sets the default command the container runs if nothing else overrides it, in this case starting the FastAPI app with uvicorn. The worker and beat services in the compose file below override this default with their own celery commands, while still building from this exact same image.
docker-compose.yml
This file defines every service in the system and how they connect to each other, letting the entire stack start with a single command.
name: fastapi-celery
services:
web:
build: .
command: uvicorn app.main:app --host 0.0.0.0 --port 8000
ports:
- "8000:8000"
env_file: .env
depends_on:
- redis
- db
worker:
build: .
command: celery -A app.core.celery_app.celery_app worker --loglevel=info -Q reports,emails --concurrency=4
env_file: .env
depends_on:
- redis
- db
beat:
build: .
command: celery -A app.core.celery_app.celery_app beat --loglevel=info
env_file: .env
depends_on:
- redis
flower:
image: mher/flower:2.0
command: celery --broker=redis://redis:6379/0 flower --port=5555
ports:
- "5555:5555"
depends_on:
- redis
redis:
image: redis:7-alpine
command: ["redis-server", "--appendonly", "yes"]
volumes:
- redis_data:/data
db:
image: postgres:16-alpine
environment:
POSTGRES_DB: ${DB_NAME}
POSTGRES_USER: ${DB_USER}
POSTGRES_PASSWORD: ${DB_PASSWORD}
volumes:
- pg_data:/var/lib/postgresql/data
volumes:
redis_data:
pg_data:
Each service explained
web builds from the Dockerfile above and runs uvicorn directly, serving the FastAPI app on port 8000, which is mapped to the same port on your host machine so you can actually reach it from a browser or a tool like curl. env_file: .env loads every variable from that file into the container's environment, which is how settings = Settings() in config.py finds DATABASE_URL and REDIS_URL once the container starts. depends_on controls startup order only, it does not wait for Redis or Postgres to be fully ready to accept connections, only for their containers to have started, which is why pool_pre_ping=True in database.py matters, it protects against the brief window where the app container is up before the database inside its own container has finished initializing.
worker builds from the identical image as web, since it needs the exact same codebase and dependencies, but overrides the default command entirely. celery -A app.core.celery_app.celery_app worker points Celery's command line tool at the celery_app object inside that module and starts a worker process against it. -Q reports,emails tells this particular worker to only consume tasks from those two named queues, matching the task_routes configuration back in celery_app.py. --concurrency=4 runs four worker processes in parallel inside this one container, each capable of executing a task independently, so four tasks can run at the same time rather than one at a time.
beat also builds from the same image, running celery ... beat instead, which starts the scheduler process responsible for pushing the scheduled tasks defined in celery_app.conf.beat_schedule onto the queue at the times configured. This service deliberately has no -Q flag and does no task execution itself, its only job is watching the clock and enqueuing work for the worker service to actually pick up and run. Exactly one beat service should ever run per deployment, running two would push every scheduled task onto the queue twice.
flower uses a prebuilt image rather than building from your own Dockerfile, since it is a standalone monitoring tool with no dependency on your application code. It connects to the same Redis broker your worker uses and exposes a web dashboard on port 5555, showing queued, running, and completed tasks along with retry counts and failure tracebacks in real time.
redis runs the official lightweight Redis image, serving as both the broker and result backend for the whole stack. command: ["redis-server", "--appendonly", "yes"] turns on append only file persistence, which means Redis writes every change to disk as it happens, so if the Redis container restarts, queued tasks that had not yet been picked up are not simply lost.
db runs PostgreSQL, with its database name, user, and password set directly through environment variables here for simplicity, though in a real deployment these would come from the same .env file rather than being hardcoded in the compose file itself. Both redis and db mount named volumes, redis_data and pg_data, declared at the bottom of the file, which is what makes their data survive a container being removed and recreated, without a volume, destroying the container destroys everything stored inside it.
Wrapping Up
Every file here does exactly one job, and the boundaries between them are what make the system debuggable later. When a report fails to generate, you know to look at tasks.py for the logic and celery_app.py for the queue it ran on. When a job gets created twice, you know to check whether service.py committed before or after dispatching the task. That discipline is a small cost upfront and a large saving every time something breaks at two in the morning.
Source Code
You can find the complete source code on GitHub. https://github.com/Obrainwave/fastapi-celery

