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, BullMQworker.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.
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.
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
}
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
- Retry Strategies & Backoff — the retry layer breakers sit above.
- Preventing Retry Storms After an Outage — what happens when the breaker closes.
- Retrying Only Transient Errors by Exception Type — deciding which failures count.
- Rate Limiting Third-Party API Calls from Workers — protecting dependencies when they are healthy.