Fixing Celery Worker Memory Leaks

Long-lived Celery pool processes run thousands of tasks each, so any memory a task fails to release accumulates until the container is OOM-killed. This guide contains the growth immediately and then finds the actual leak, as part of Celery Architecture & Configuration in Backend Frameworks & Worker Scaling.

Problem Statement

A reporting service's Celery workers start at about 250 MB per pool process and grow steadily; after 8–12 hours a pod with 8 processes exceeds its 4 GB limit and is OOM-killed, interrupting whatever was running. Restarts are frequent enough that tasks occasionally run twice, and a scheduled restart cron hides the problem without explaining it. You want memory bounded right away without waiting for a fix, the growth traced to specific tasks and allocations, and a durable fix with a regression check.

Prerequisites

  • Celery 5.3+ with the prefork pool (the recycling settings apply to pool children).
  • Python 3.10+ for tracemalloc; optionally memray or objgraph for deeper analysis.
  • Metrics for per-process memory (node or cAdvisor metrics per container, plus per-child RSS if possible).
  • A staging environment where a single worker can be run against replayed tasks.

Step 1 — Contain the Growth with Child Recycling

Before hunting the leak, stop it from taking pods down. Celery can replace a pool child after it has run a number of tasks or exceeded a memory threshold; the new child starts clean.

# celeryconfig.py
worker_max_tasks_per_child = 500          # recycle after 500 tasks
worker_max_memory_per_child = 600_000     # KB: recycle after the child exceeds ~600 MB RSS

worker_max_memory_per_child is checked after a task finishes: a child exceeding the threshold finishes its current task and is then replaced. It prevents slow leaks from growing forever but cannot stop a single task that allocates gigabytes. Size the threshold so concurrency × threshold plus the main process fits comfortably under the container limit.

Recycling bounds a slow leak Without recycling, each pool process grows from 250 megabytes by about 40 megabytes an hour, and after about ten hours eight processes exceed the four-gigabyte pod limit and the pod is killed. With worker_max_memory_per_child set to 600 megabytes, each child is replaced after the task that pushes it past the threshold, so memory saws between 250 and 600 megabytes indefinitely. Per-child RSS over 12 hours limit / 8 no recycling: OOM kill recycled at 600 MB 0 h 12 h

Step 2 — Attribute Growth to Task Types

A leak is usually caused by one task type. Record RSS before and after each task per child, and aggregate the difference by task name.

import os, resource, psutil
from celery.signals import task_prerun, task_postrun

_proc = psutil.Process(os.getpid())
_before: dict[str, int] = {}

@task_prerun.connect
def rss_before(task_id, task, **_):
    _before[task_id] = _proc.memory_info().rss

@task_postrun.connect
def rss_after(task_id, task, **_):
    start = _before.pop(task_id, None)
    if start is not None:
        delta = _proc.memory_info().rss - start
        TASK_RSS_DELTA.labels(task=task.name).observe(max(delta, 0))
        CHILD_RSS.labels(pid=str(os.getpid())).set(_proc.memory_info().rss)
# Average RSS retained per execution, by task: the leaking task stands out
sum by (task) (rate(celery_task_rss_delta_bytes_sum[1h]))
  / sum by (task) (rate(celery_task_rss_delta_bytes_count[1h]))

Individual deltas are noisy (the allocator keeps freed memory), but averaged over thousands of executions, a task that retains memory shows a consistently positive value. In the scenario, build_customer_report retained about 180 KB per run; every other task averaged near zero.

Step 3 — Find the Allocation with tracemalloc

Run a single worker with one process, replay the suspect task a few hundred times, and compare tracemalloc snapshots.

# scripts/leak_hunt.py
import tracemalloc, gc
from tasks import build_customer_report

tracemalloc.start(25)
for i in range(50):                       # warm up caches and imports first
    build_customer_report.run(sample_ids[i % len(sample_ids)])
gc.collect()
before = tracemalloc.take_snapshot()
for i in range(500):
    build_customer_report.run(sample_ids[i % len(sample_ids)])
gc.collect()
after = tracemalloc.take_snapshot()
for stat in after.compare_to(before, "traceback")[:10]:
    print(stat)
    for line in stat.traceback.format()[-6:]:
        print("   ", line)

The top entries show where retained memory was allocated, with a stack. Warming up first matters: the first calls legitimately fill caches and import modules, which is growth but not a leak. For leaks in C extensions, which tracemalloc does not see, use memray run --native on the same script.

Snapshot, repeat, compare First the suspect task runs fifty times to fill caches, then a tracemalloc snapshot is taken. The task runs five hundred more times and a second snapshot is taken after garbage collection. Comparing the two by traceback lists the allocation sites whose memory grew, with the code path that created them. Isolate the leak in one process warm up 50 runs snapshot A repeat 500 runs snapshot B, compare top growth by traceback Warm-up separates one-time cache fills from growth that repeats on every run.

Step 4 — Check the Usual Suspects

Most worker leaks come from a short list of causes, all of which follow from the process living far longer than any one task:

# 1. Unbounded module-level caches
_cache = {}                                       # grows with every distinct customer
def get_template(customer_id):
    if customer_id not in _cache:
        _cache[customer_id] = load_template(customer_id)
    return _cache[customer_id]
# fix: functools.lru_cache(maxsize=256) or a TTL cache

# 2. Accumulating lists on long-lived objects
class ReportBuilder:                               # created once at import
    rows = []                                      # class attribute shared across calls
# fix: create per-task state inside the task

# 3. ORM sessions or query caches that are never cleared
#    e.g. Django DEBUG=True stores every query in connection.queries
from django.db import reset_queries               # call per task, and never run DEBUG in workers

# 4. Global event handlers / signal receivers registered per call
#    each task run adds another handler that references large objects

In the scenario, the culprit was pattern 1: a template cache keyed by customer that grew with every new customer processed. An LRU bound fixed it: memory now plateaus at the size of the 256 most recently used templates, instead of growing with the total number of customers ever processed by that child.

Bounded caches plateau; unbounded ones leak A module-level dictionary cache keyed by customer adds a template for every new customer a child process sees, so its size grows in proportion to tasks run. An LRU cache with a maximum of 256 entries grows until it holds 256 templates and then evicts the least recently used, so memory flattens. Cache size vs distinct customers processed dict: grows forever lru_cache(256): plateaus 0 customers seen by one child

Django's DEBUG=True (pattern 3) is common enough to check first in any Django project.

Step 5 — Set Kubernetes Limits That Match the Configuration

Container limits and Celery settings must agree, or the OOM killer acts before Celery's recycling can.

resources:
  requests: { cpu: "2", memory: "3Gi" }
  limits:   { memory: "5Gi" }        # 8 children x 600 MB + main process + headroom
env:
  - { name: CELERY_WORKER_MAX_MEMORY_PER_CHILD, value: "600000" }

When the kernel OOM-kills a child, Celery's main process sees WorkerLostError; with task_reject_on_worker_lost the task is requeued, as described in Celery acks_late and worker crash safety. When it kills the main process, the pod restarts. Recycling at a threshold well below the limit keeps both events rare.

Step 6 — Add a Regression Check

After the fix, keep a test that fails if the task starts retaining memory again. It does not need to be precise — only to catch a return of steady growth.

def test_report_task_does_not_retain_memory():
    for i in range(50):
        build_customer_report.run(fixture_ids[i % 10])
    gc.collect()
    tracemalloc.start()
    base = tracemalloc.get_traced_memory()[0]
    for i in range(300):
        build_customer_report.run(fixture_ids[i % 10] + i)    # distinct ids: exercises caches
    gc.collect()
    grown = tracemalloc.get_traced_memory()[0] - base
    assert grown < 5 * 1024 * 1024, f"retained {grown / 1e6:.1f} MB over 300 runs"

Using distinct ids is what makes the test catch the unbounded-cache class of leak; repeating the same input would hide it.

Verification

After deploying the fix, per-child RSS should plateau instead of rising, recycling events from worker_max_memory_per_child should become rare, and OOM kills (kube_pod_container_status_last_terminated_reason{reason="OOMKilled"}) should stop.

max by (pod) (container_memory_working_set_bytes{container="celery-worker"})
increase(kube_pod_container_status_restarts_total{container="celery-worker"}[24h])

Keep the recycling settings even after the fix; they are cheap insurance against the next leak.

Gotchas & Edge Cases

RSS never shrinks. Python's allocator and glibc keep freed memory for reuse, so RSS often stays high after a spike without any leak. Look for steady growth across many tasks, not a high plateau.

Very large single tasks. A task that loads a 2 GB file will exceed any per-child threshold on its own. Recycling does not help; stream the data instead.

Recycling cost. Each new child re-imports the app (seconds for large Django projects). worker_max_tasks_per_child too low wastes CPU; favour the memory threshold.

Memory growth in the main process. The worker's main process also runs code — signal handlers, event dispatchers, custom bootsteps. A leak there is not fixed by child recycling. Watch the main process's RSS separately from the children's.

Fork and copy-on-write. Children share the parent's memory pages until they write to them; touching large imported objects (even by reference counting) gradually copies them. gc.freeze() in the parent before forking reduces this effect for large Django or data-science imports.

Other pools. With gevent or threads, there are no children to recycle; the whole worker must be restarted, so leaks matter more.

FAQ

Is a nightly worker restart an acceptable fix? It is containment, like recycling, but coarser — and it interrupts running tasks. Prefer child recycling, and still find the leak.

Does memory growth always mean a leak? No. Caches warming up and allocator behaviour both grow RSS. A leak is growth proportional to the number of tasks run, without a plateau.

How do I profile in production without slowing workers down? tracemalloc adds noticeable overhead, so keep it for staging replays. In production, the per-task RSS delta from Step 2 is cheap enough to leave on permanently, and py-spy dump or memray attach can inspect a single live process briefly when you need a snapshot of a specific leaking child.

Should I lower concurrency to avoid OOM kills? Only as a temporary measure. Fewer processes per pod means each can grow further before the pod limit, which delays the kill without fixing the growth — and it lowers throughput. Recycling plus a fix is the durable answer.

Do the same tools work for Sidekiq or RQ? The approach does: attribute growth per job, isolate with heap snapshots (memory_profiler, derailed_benchmarks for Ruby), and bound process lifetime. See diagnosing Sidekiq memory bloat.

Related