Circuit Breakers for Worker Dependencies

Retries handle a dependency that fails occasionally; a circuit breaker handles one that is down. This guide adds breakers to queue workers so a failing downstream is left alone to recover instead of being hammered by every retry, as part of Retry Strategies & Backoff in Queue Fundamentals & Architecture.

Problem Statement

A fulfilment service's workers call a warehouse API for every order job. When the warehouse API went down for 35 minutes, 40 workers kept pulling jobs, each call waited for a 10-second timeout, and each failure scheduled a retry. Worker slots were tied up waiting on timeouts, unrelated jobs on the same workers stalled, the warehouse API's recovery was slowed by a flood of connection attempts, and 3,000 jobs exhausted their retries during the outage and landed in the dead-letter queue — for an outage that needed no action from anyone. You want workers to detect the outage within seconds, stop calling the API fleet-wide, stop consuming the jobs that depend on it without spending retries, probe for recovery, and resume automatically.

Prerequisites

  • Workers that call identifiable dependencies (per host or per API) and can check a breaker before calling.
  • A shared store for breaker state across workers — Redis here — because a per-process breaker needs every process to rediscover the outage.
  • Timeouts on every outbound call (a breaker cannot open on calls that never return).
  • A way to pause consumption of a queue (Celery cancel_consumer, BullMQ worker.pause, Sidekiq quiet, or simply deferring jobs).

Step 1 — Understand the Three States

A breaker is a small state machine per dependency:

  • Closed: calls go through; failures are counted in a rolling window.
  • Open: calls fail immediately without touching the dependency, for a cool-down period.
  • Half-open: after the cool-down, a limited number of trial calls are allowed; success closes the breaker, failure reopens it.
closed --(failure rate >= 50% over last 20 calls, min 10 calls)--> open
open   --(after 30 s cool-down)------------------------------------> half-open
half-open --(3 trial calls succeed)-------------------------------> closed
half-open --(any trial call fails)---------------------------------> open (cool-down doubles, max 5 min)

Using a failure rate with a minimum call count, rather than N consecutive failures, avoids opening on a couple of unlucky requests during low traffic and avoids staying closed when a few successes are interleaved with mostly failures.

Closed, open, half-open In the closed state calls pass through and failures are counted. When the failure rate over the recent window reaches fifty percent with at least ten calls, the breaker opens and calls fail fast for thirty seconds. It then becomes half-open and allows three trial calls. If they succeed it closes; if any fails it reopens with a doubled cool-down. Breaker state machine closed calls pass, count failures open fail fast, cool down half-open 3 trial calls rate >= 50% after 30 s trials succeed: back to closed; any trial fails: back to open, cool-down doubles

Step 2 — Keep Breaker State in Redis So the Fleet Agrees

A per-process breaker means 40 processes each need 10 failures to learn the dependency is down — 400 wasted calls and 400 timeout waits. A shared breaker opens once for everyone.

# breaker.py — shared, rate-based circuit breaker
import time, redis
r = redis.Redis.from_url(REDIS_URL, decode_responses=True)

class Breaker:
    def __init__(self, name, window=20, min_calls=10, threshold=0.5, cooldown=30, max_cooldown=300, trials=3):
        self.k = f"cb:{name}"
        self.window, self.min_calls, self.threshold = window, min_calls, threshold
        self.cooldown, self.max_cooldown, self.trials = cooldown, max_cooldown, trials

    def state(self) -> str:
        opened_until = float(r.hget(self.k, "open_until") or 0)
        if opened_until > time.time():
            return "open"
        return "half_open" if r.hget(self.k, "half_open") == "1" else "closed"

    def allow(self) -> bool:
        s = self.state()
        if s == "closed":
            return True
        if s == "open":
            return False
        return r.hincrby(self.k, "trials_started", 1) <= self.trials     # limit trial calls

    def record(self, ok: bool) -> None:
        pipe = r.pipeline()
        pipe.lpush(f"{self.k}:results", "1" if ok else "0")
        pipe.ltrim(f"{self.k}:results", 0, self.window - 1)
        pipe.lrange(f"{self.k}:results", 0, -1)
        results = pipe.execute()[2]
        if self.state() == "half_open":
            if not ok:
                return self._open(double=True)
            if r.hincrby(self.k, "trials_ok", 1) >= self.trials:
                r.delete(self.k, f"{self.k}:results")                     # fully closed, fresh window
            return
        failures = results.count("0")
        if len(results) >= self.min_calls and failures / len(results) >= self.threshold:
            self._open(double=False)

    def _open(self, double: bool) -> None:
        cd = min(float(r.hget(self.k, "cooldown") or self.cooldown) * (2 if double else 1), self.max_cooldown)
        r.hset(self.k, mapping={"open_until": time.time() + cd, "cooldown": cd, "half_open": "1",
                                "trials_started": 0, "trials_ok": 0})

After the cool-down, state() reports half-open (because half_open is set and open_until has passed), a bounded number of workers make trial calls, and the breaker either clears or reopens with a longer cool-down. The small races between workers updating the list are acceptable for a breaker; a Lua script can make the update atomic if you need precision.

Step 3 — Wrap Dependency Calls and Fail Fast

Every call to the dependency goes through the breaker. When open, the call raises immediately — no connection attempt, no timeout wait.

class DependencyUnavailable(Exception):
    """Raised without calling the dependency when its breaker is open."""

warehouse_breaker = Breaker("warehouse-api")

def reserve_stock(order):
    if not warehouse_breaker.allow():
        raise DependencyUnavailable("warehouse-api")
    try:
        resp = warehouse.post("/reservations", json=order.lines, timeout=5)
        resp.raise_for_status()
    except (httpx.TimeoutException, httpx.ConnectError, httpx.HTTPStatusError) as exc:
        if isinstance(exc, httpx.HTTPStatusError) and exc.response.status_code < 500:
            warehouse_breaker.record(ok=True)          # 4xx: the service is up, the request is bad
            raise
        warehouse_breaker.record(ok=False)
        raise
    warehouse_breaker.record(ok=True)
    return resp.json()

Only failures that indicate the dependency is unhealthy (timeouts, connection errors, 5xx) count against the breaker. A 400 or 404 means the service answered; counting those would let bad input open the breaker for everyone.

Step 4 — Defer Jobs Without Burning Retries While Open

Raising DependencyUnavailable inside a job must not count as a failed attempt, or the outage would still drain retry budgets into the DLQ. Defer the job until about when the breaker will half-open.

@app.task(bind=True, acks_late=True, max_retries=6, autoretry_for=(httpx.TransportError,),
          retry_backoff=True)
def fulfil_order(self, order_id):
    try:
        reserve_stock(load_order(order_id))
    except DependencyUnavailable:
        # not a failure of this job: requeue without counting an attempt
        raise self.retry(countdown=30 + random.uniform(0, 15), max_retries=None)

For even less churn, stop consuming the dependent queue while the breaker is open: a small watcher cancels consumption of fulfilment when the breaker opens and resumes it when it closes, so jobs simply wait in the broker instead of cycling through workers.

def breaker_watcher():
    consuming = True
    while True:
        s = warehouse_breaker.state()
        if s == "open" and consuming:
            app.control.cancel_consumer("fulfilment")       # all workers stop taking these jobs
            consuming = False
        elif s != "open" and not consuming:
            app.control.add_consumer("fulfilment")
            consuming = True
        time.sleep(2)

Other queues on the same workers keep flowing, which removes the "unrelated jobs stall" symptom. The equivalent in BullMQ is worker.pause() on the dependent queue's workers; in Sidekiq, a separate process for the dependent queue that is quieted. Retry budgets in general are covered in setting retry budgets and max attempts.

Outage with and without a breaker Without a breaker, every worker keeps calling the dead API, each call waiting ten seconds for a timeout, and 3,000 jobs exhaust their retries into the dead-letter queue. With a shared breaker, it opens after about ten failed calls, consumption of the fulfilment queue pauses, jobs wait in the broker, half-open trials detect recovery after 35 minutes, and the backlog drains with no dead letters. 35-minute warehouse API outage no breaker timeouts, retries, slots blocked 3,000 in DLQ breaker open: consumption paused, jobs wait closed: backlog drains The breaker opened after ~10 failures; half-open trials found recovery within one cool-down. Zero jobs dead-lettered, zero manual replays.

Step 5 — Use One Breaker per Dependency, Not per Job Type

Scope breakers to what fails together: one per external API, per database, or per region of a service. A single global breaker would stop all work when one dependency fails; one breaker per job type would make each job type rediscover the same outage.

BREAKERS = {
    "warehouse-api": Breaker("warehouse-api"),
    "payments-api": Breaker("payments-api", threshold=0.3, cooldown=15),   # tighter for money paths
    "geocoder": Breaker("geocoder", min_calls=30),                        # noisy API: needs more samples
}
Scope the breaker to what fails together A single global breaker opens for everything when one API fails. Per-job-type breakers make three job types that all call the warehouse API each rediscover the same outage. Per-dependency breakers let the warehouse-api breaker open while payments and geocoding keep working. Warehouse API down: which breakers open? one global breaker payments stop too per job type 3 breakers learn it 3 times per dependency only warehouse-api opens The right-hand scope matches the failure domain; the others are too broad or too narrow.

Where a dependency has several independent endpoints or regions, break per endpoint so an outage in one region does not stop calls to the others.

Step 6 — Observe Breaker Transitions

A breaker that opens is an incident signal in its own right, often earlier and clearer than error rates.

BREAKER_STATE = Gauge("circuit_breaker_open", "1 when open or half-open", ["dependency"])
BREAKER_TRANSITIONS = Counter("circuit_breaker_transitions_total", "", ["dependency", "to"])
# Page when a critical dependency's breaker has been open for 5 minutes
max by (dependency) (circuit_breaker_open{dependency=~"payments-api|warehouse-api"}) == 1
# (with for: 5m in the alert rule)

# Flapping: frequent transitions mean thresholds are too sensitive or the dependency is degraded
sum by (dependency) (increase(circuit_breaker_transitions_total[15m])) > 6

Show breaker state on the queue dashboard next to backlog: a growing backlog with an open breaker is expected and self-healing; a growing backlog with a closed breaker needs attention elsewhere.

Verification

def test_breaker_opens_and_recovers(fake_clock, flaky_warehouse):
    b = Breaker("test", window=10, min_calls=5, cooldown=30, trials=2)
    flaky_warehouse.down()
    for _ in range(5):
        assert b.allow(); b.record(ok=False)
    assert b.state() == "open" and not b.allow()
    fake_clock.advance(31)
    flaky_warehouse.up()
    assert b.allow(); b.record(ok=True)
    assert b.allow(); b.record(ok=True)
    assert b.state() == "closed"

In staging, block the dependency with a network policy under load and confirm: the breaker opens within seconds, the dependent queue's consumption stops, no jobs reach the DLQ, and processing resumes within one cool-down of unblocking.

Gotchas & Edge Cases

No timeouts, no breaker. A call without a timeout can hang indefinitely and never record a failure. Set connect and read timeouts on every client.

Breakers masking partial failures. A dependency that fails only for some inputs (one bad region, one large tenant) can open the breaker for everyone. Break on dependency health, and handle input-specific errors in the job.

Redis unavailable. Decide the fallback: failing open (calls allowed) keeps work flowing without protection; failing closed stops everything. Most teams fail open with a local per-process breaker as backup.

Half-open stampede. Without a trial limit, every worker tries the moment the cool-down ends. Keep trials small.

FAQ

Aren't retries with backoff enough? Backoff spaces retries for one job; it does not stop 40 workers from each discovering the outage independently, nor stop new jobs from trying. A breaker adds the fleet-wide "stop calling" decision.

Should the breaker live in the worker or a service mesh? A mesh (Envoy outlier detection, Istio) protects the dependency but cannot pause queue consumption or keep retry budgets intact. Worker-level breakers can; the two complement each other.

How long should the cool-down be? Start around 15–30 seconds and double on repeated failures up to a few minutes. Short enough to notice recovery quickly, long enough not to probe a struggling service constantly.

Related