Designing Circuit Breakers for Webhook Ingress: Preventing Database Lock Contention During Downstream Outages
Learn how to implement dynamic ingestion circuit breakers, manage DLQ pressure, and optimize backoff policies to protect downstream infrastructure during outages.

Designing Circuit Breakers for Webhook Ingress: Preventing Database Lock Contention During Downstream Outages
Executive Summary
Webhook pipelines take in asynchronous event payloads such as payment confirmations, order updates and shipment milestones, and persist them into internal systems. When a downstream dependency such as a relational database slows down, naive retry logic can turn a local slowdown into a system-wide outage.
The mechanism is simple. Retried messages and fresh webhooks compete for the same rows, workers block on locks, and connection pools drain. Healthy requests that touch unrelated data then fail because they cannot get a connection, and those failures generate more retries.
This guide describes a dynamic ingestion circuit breaker that sits between your queue and your database. It trips on queue depth and in-flight metrics, dead-letter queue (DLQ) pressure, and worker-reported lock waits. It pairs with database-side lock timeouts and retries that use full jitter. While the breaker is open, webhooks wait in a durable queue instead of hammering a struggling database. With idempotent consumers, this gives you at-least-once processing across the outage, bounded by the queue's retention period.
The thresholds in this article are illustrative starting points, not universal constants. Tune them against your own baselines.
1. The Anatomy of a Cascading Failure
1.1 The Naive Retry Antipattern
A typical webhook architecture has an HTTP endpoint, a durable queue (for example Amazon SQS, RabbitMQ or Apache Kafka), a fleet of consumer workers, and a relational database such as PostgreSQL.
+----------------+ +--------------------+ +---------------+ +--------------+
| Webhook Issuer |-->| Ingress API Gateway|-->| Ingress Queue |-->| Worker Fleet |
+----------------+ +--------------------+ +---------------+ +------+-------+
|
v
+----------------+
| Database/Store |
+----------------+
The endpoint should do as little as possible: verify the signature, enqueue the payload, and return a 2xx quickly. This matters because issuers enforce tight response deadlines and have their own retry policies:
- Shopify treats a delivery as failed if the app does not respond with a 2xx status within five seconds. It retries failed deliveries up to eight times over four hours, with the interval increasing each time, and removes the webhook subscription if failures persist.
- Stripe, in live mode, retries delivery of an event for up to three days with exponential backoff. In test mode it retries three times over a few hours.
Decoupling ingestion from processing therefore protects you in two ways. A database problem does not turn into HTTP 5xx responses to the issuer, and it does not eat into the issuer's retry budget.
The failure starts inside the worker fleet. When a processing error occurs, such as an HTTP timeout or a database lock wait, naive workers typically:
- Push the message back onto the main queue or a retry topic.
- Retry at fixed intervals, or with uncapped exponential backoff and no jitter.
- Keep pulling new messages at full concurrency regardless of how many are failing.
1.2 The Lock Contention Spiral
Suppose the root cause is database-side: a long-running transaction, an unindexed update, or a slow disk. Aggressive retries make it worse. Consider an order-status webhook consumer:
- Transaction A updates
orderswhereid = 1042, taking a row-level lock that is held until commit or rollback. - A storage slowdown stretches Transaction A from about 10 ms to several seconds.
- Dozens of webhooks for the same order arrive and each worker opens a transaction. All of them block waiting for that row lock. In PostgreSQL, such sessions show a
Lockwait event inpg_stat_activity. - Blocked workers hold their connections, so the pool (HikariCP,
pg.Pool, or a PgBouncer limit) is exhausted. - Webhooks for unrelated orders (
1043,1044) cannot get a connection, fail, and re-enter the retry flow. - The retry flow multiplies lock requests and connection attempts until the database saturates.
Incoming Webhook Burst
|
v
+-------------------------------+
| Active Workers (High Threads) |
+-------------------------------+
|
+--> [Lock request: order 1042] --> (Executing - slow)
+--> [Lock request: order 1042] --> (Waiting on lock...)
+--> [Lock request: order 1042] --> (Waiting on lock...)
|
v
+-------------------------------+
| DB connection pool exhausted |
+-------------------------------+
|
v
Failures across ALL webhook handlers, including unrelated entities
This follows a pattern the AWS Builders' Library describes for retries in general: retries add load to a dependency that is already struggling, so they need limits and need to be coordinated with other protective mechanisms.
2. Circuit Breakers vs. Backoff Policies
Retries with backoff and circuit breaking solve different problems, and a robust system uses both. The circuit breaker pattern, as popularized by Michael Nygard's Release It! and described by Martin Fowler, wraps a protected call and stops making that call once failures cross a threshold.
| Dimension | Exponential backoff with jitter | Circuit breaker |
|---|---|---|
| Scope | Individual message or call | The whole path to a failing resource |
| Mechanism | Delays re-execution of one payload | Stops or throttles all calls to the resource |
| Best for | Transient network glitches, brief races | Sustained outages, lock exhaustion, overload |
| Load on the dependency | Still sends work, just spread out | Sheds load while open |
2.1 Why Full Jitter Matters
Capped exponential backoff computes the delay as:
twait=min(tmax,tbase×2attempt)t_{\text{wait}} = \min(t_{\text{max}}, t_{\text{base}} \times 2^{\text{attempt}})twait=min(tmax,tbase×2attempt)
If many workers fail at the same moment, they all retry at the same future instants (2 s, 4 s, 8 s and so on) and produce synchronized bursts, the thundering herd. Full jitter randomizes the whole delay:
tsleep=random(0,min(tmax,tbase×2attempt))t_{\text{sleep}} = \text{random}(0, \min(t_{\text{max}}, t_{\text{base}} \times 2^{\text{attempt}}))tsleep=random(0,min(tmax,tbase×2attempt))
AWS's Architecture Blog post "Exponential Backoff and Jitter" simulated contending clients and showed that adding jitter substantially reduces total calls and completion time compared with backoff alone. Stripe's guidance for handling rate limits similarly recommends exponential backoff with some randomness to avoid a thundering herd.
2.2 Why Jitter Alone Is Not Enough
Jitter spreads retries out, but it does not reduce the total work offered to the database. If the arrival rate of new webhooks exceeds what a degraded database can process, retries keep adding to the load. A dynamic ingestion circuit breaker operates above the message level. It watches system-level health signals and tells workers to pause polling or cut concurrency until the dependency recovers.
3. Architecture of a Dynamic Ingestion Circuit Breaker
Rather than judging each message in isolation, the controller evaluates three signals: in-flight vs. visible messages, DLQ pressure, and worker heartbeats.
+------------------------+
| Breaker Controller |
| (control plane) |
+-----------+------------+
^
+--------------------------+--------------------------+
| | |
[Signal 1: DLQ pressure] [Signal 2: In-flight ratio] [Signal 3: Heartbeats]
| | |
+------+-------+ +-------+-------+ +-------+-------+
| Dead-letter | | Ingress | | Worker fleet |
| queue | | queue | | |
+--------------+ +---------------+ +---------------+
3.1 State Machine
The breaker has three states:
- CLOSED (normal): workers pull at full capacity, and failure signals stay below thresholds.
- OPEN (paused or throttled): triggered when error rates, lock-timeout rates or DLQ velocity cross thresholds. Workers stop polling the ingress queue, or sharply reduce concurrency and batch size. Incoming webhooks accumulate in the durable queue.
- HALF-OPEN (probing): after a cool-down period, a small number of probe messages are allowed through. If they succeed, the breaker returns to CLOSED. If they fail, it returns to OPEN and the cool-down restarts.
Recovery should be gradual. Releasing a large backlog at full concurrency the instant the breaker closes can re-trip it, so ramp concurrency back up in steps.
4. Signals for Tripping the Breaker
4.1 DLQ Pressure
A DLQ receives messages that have exhausted their retries. In a database incident, DLQ arrival rate spikes. Define DLQ arrival velocity over a rolling window:
Vdlq=ΔNdlqΔtV_{\text{dlq}} = \frac{\Delta N_{\text{dlq}}}{\Delta t}Vdlq=ΔtΔNdlq
Trip the breaker when VdlqV_{\text{dlq}} Vdlq exceeds a threshold, or when the failure ratio exceeds a limit (for example more than 15% of processed messages over one minute). Both values are starting points to tune against your normal error rate.
On Amazon SQS, DLQ routing is built in. A redrive policy moves a message to the DLQ after it has been received maxReceiveCount times without being deleted. Kafka has no native DLQ, so teams publish failed records to a dedicated topic.
4.2 In-Flight Ratio
When workers block on locks, messages stay in flight (received but neither deleted nor completed). On SQS you can read this from CloudWatch: ApproximateNumberOfMessagesNotVisible counts in-flight messages and ApproximateNumberOfMessagesVisible counts messages waiting to be received. Define:
Rinflight=MinflightMinflight+MvisibleR_{\text{inflight}} = \frac{M_{\text{inflight}}}{M_{\text{inflight}} + M_{\text{visible}}}Rinflight=Minflight+MvisibleMinflight
A high ratio combined with falling throughput suggests workers are stuck. A reasonable starting rule is to open the breaker when Rinflight>0.85R_{\text{inflight}} > 0.85 Rinflight>0.85 and throughput drops below 20% of baseline. Treat the ratio with care on its own, because a legitimately drained queue also produces a high ratio. Always pair it with throughput.
Also be aware of broker limits. SQS standard queues allow approximately 120,000 in-flight messages and FIFO queues allow 20,000. Hitting the limit stops you from receiving more messages, so a stuck fleet can starve the queue.
4.3 Worker Heartbeats and Lock Detection
Each worker can publish a short-lived heartbeat to a shared store such as Redis:
{
"worker_id": "worker-node-8f4a",
"thread_id": 14,
"state": "DB_LOCK_WAIT",
"lock_target": "orders:1042",
"wait_duration_ms": 4200,
"timestamp": 1728000000
}
If, say, more than 30% of active workers report lock waits above 2,000 ms, the controller trips the breaker before the pool runs dry.
You can complement worker-reported state with the database's own view. In PostgreSQL, pg_stat_activity shows sessions waiting on locks (wait_event_type = 'Lock'), pg_locks shows lock requests, and pg_blocking_pids(pid) identifies which sessions are blocking others. Setting log_lock_waits makes PostgreSQL log waits that exceed deadlock_timeout.
5. Database-Side Guardrails
The breaker works best when each worker's database session also fails fast.
lock_timeout: aborts a statement that waits longer than the configured time to acquire a lock. The error has SQLSTATE55P03(lock_not_available). Per-role or per-session configuration is usually safer than a global setting.statement_timeout: note that lock wait time counts toward the statement timeout, so astatement_timeoutshorter thanlock_timeoutmakes the lock timeout moot.idle_in_transaction_session_timeout: terminates sessions that open a transaction and then stall, which would otherwise hold locks indefinitely.- Bounded pools: cap pool size and use a short connection-acquisition timeout so that saturation surfaces as a fast, classifiable error.
- Short transactions and consistent lock ordering: keep transactions small, and acquire locks in the same order everywhere to reduce deadlocks and wait chains.
6. End-to-End Node.js Implementation
The example below combines a Redis-backed breaker, full-jitter retries, and lock-timeout classification. It is a teaching sketch: production code would use the broker's own acknowledgement semantics and DLQ, structured logging, and metrics.
Three design points are worth noting:
- Workers must not poll while the breaker is open. Receiving a message and then releasing it is wasteful, and on SQS each receive increments the message's receive count toward
maxReceiveCount. - In HALF_OPEN, a Redis
SET NXlock admits exactly one probe at a time instead of letting every worker through. - The consumer must be idempotent. At-least-once delivery means duplicates will occur.
/**
* Webhook consumer with a Redis-backed circuit breaker.
* Stack: Node.js (ESM), ioredis, node-postgres (pg)
*/
import Redis from 'ioredis';
import pg from 'pg';
const { Pool } = pg;
const CONFIG = {
REDIS_URL: process.env.REDIS_URL || 'redis://127.0.0.1:6379',
POSTGRES_URL: process.env.POSTGRES_URL || 'postgres://user:pass@localhost:5432/webhooks',
STATE_KEY: 'cb:webhook_ingress:state',
FAILURE_COUNTER_KEY: 'cb:webhook_ingress:failures',
PROBE_KEY: 'cb:webhook_ingress:probe',
FAILURE_THRESHOLD: 10, // failures allowed per window
WINDOW_SECONDS: 30, // rolling-ish window (fixed window via key expiry)
COOLDOWN_MS: 15000, // time spent OPEN before probing
PROBE_TTL_MS: 10000, // max time a probe may hold the probe lock
MAX_RETRIES: 3,
BASE_BACKOFF_MS: 100,
MAX_BACKOFF_MS: 3000,
};
const redis = new Redis(CONFIG.REDIS_URL);
const dbPool = new Pool({
connectionString: CONFIG.POSTGRES_URL,
max: 20, // bounded pool
idleTimeoutMillis: 30000,
connectionTimeoutMillis: 2000, // fail fast when the pool is saturated
});
class CircuitBreaker {
static CLOSED = 'CLOSED';
static OPEN = 'OPEN';
static HALF_OPEN = 'HALF_OPEN';
/** OPEN becomes HALF_OPEN implicitly once the cool-down has elapsed. */
static async getState() {
const raw = await redis.get(CONFIG.STATE_KEY);
if (!raw) return this.CLOSED;
const data = JSON.parse(raw);
if (data.state === this.OPEN && Date.now() - data.openedAt >= CONFIG.COOLDOWN_MS) {
return this.HALF_OPEN;
}
return data.state;
}
/** In HALF_OPEN, only the worker that wins this lock may send a probe. */
static async tryAcquireProbe() {
const res = await redis.set(CONFIG.PROBE_KEY, '1', 'PX', CONFIG.PROBE_TTL_MS, 'NX');
return res === 'OK';
}
static async open(reason) {
console.warn(`[BREAKER] OPEN: ${reason}`);
await redis.set(CONFIG.STATE_KEY, JSON.stringify({
state: this.OPEN,
openedAt: Date.now(),
reason,
}));
await redis.del(CONFIG.PROBE_KEY);
}
static async close() {
console.info('[BREAKER] Probe succeeded. Resetting to CLOSED.');
await redis.del(CONFIG.STATE_KEY, CONFIG.FAILURE_COUNTER_KEY, CONFIG.PROBE_KEY);
}
static async recordFailure(kind) {
const count = await redis.incr(CONFIG.FAILURE_COUNTER_KEY);
if (count === 1) {
await redis.expire(CONFIG.FAILURE_COUNTER_KEY, CONFIG.WINDOW_SECONDS);
}
console.error(`[FAILURE] ${count}/${CONFIG.FAILURE_THRESHOLD} (${kind})`);
if (count >= CONFIG.FAILURE_THRESHOLD) {
await this.open(`${count} failures within ${CONFIG.WINDOW_SECONDS}s`);
}
}
}
/** Full jitter: random(0, min(cap, base * 2^attempt)) */
export function fullJitterBackoff(attempt) {
const ceiling = Math.min(CONFIG.MAX_BACKOFF_MS, CONFIG.BASE_BACKOFF_MS * 2 ** attempt);
return Math.floor(Math.random() * ceiling);
}
const sleep = (ms) => new Promise((resolve) => setTimeout(resolve, ms));
/** Idempotent upsert guarded by a session-local lock timeout. */
async function persistWebhook(payload) {
const client = await dbPool.connect();
try {
await client.query('BEGIN');
// SET LOCAL lasts until the end of this transaction.
await client.query("SET LOCAL lock_timeout = '1000ms'");
await client.query(
`INSERT INTO webhook_events (event_id, entity_id, payload, processed_at)
VALUES ($1, $2, $3, NOW())
ON CONFLICT (event_id)
DO UPDATE SET payload = EXCLUDED.payload, processed_at = NOW()`,
[payload.id, payload.entityId, JSON.stringify(payload)]
);
await client.query('COMMIT');
} catch (err) {
await client.query('ROLLBACK').catch(() => {});
if (err.code === '55P03') { // lock_not_available (lock_timeout expired)
const e = new Error(`Lock timeout on entity ${payload.entityId}`);
e.kind = 'LOCK_TIMEOUT';
throw e;
}
if (/timeout exceeded when trying to connect/i.test(err.message)) {
err.kind = 'POOL_TIMEOUT'; // pg pool could not hand out a connection in time
}
throw err;
} finally {
client.release();
}
}
async function sendToDeadLetterQueue(payload, error) {
// Demo only: in production use the broker's native DLQ or a dedicated topic.
await redis.rpush('queue:webhook_dlq', JSON.stringify({
payload,
error: error.message,
failedAt: new Date().toISOString(),
}));
}
/**
* Returns a status string. The caller decides whether to ack the message:
* ack on SUCCESS / DEAD_LETTERED; do NOT receive/ack when BREAKER_OPEN.
*/
export async function consumeWebhook(payload) {
const state = await CircuitBreaker.getState();
if (state === CircuitBreaker.OPEN) return 'BREAKER_OPEN';
if (state === CircuitBreaker.HALF_OPEN) {
if (!(await CircuitBreaker.tryAcquireProbe())) return 'BREAKER_OPEN';
try {
await persistWebhook(payload); // single attempt, no retries
await CircuitBreaker.close();
return 'SUCCESS';
} catch (err) {
await CircuitBreaker.open(`Probe failed: ${err.message}`);
return 'BREAKER_OPEN';
}
}
for (let attempt = 1; attempt <= CONFIG.MAX_RETRIES; attempt++) {
try {
await persistWebhook(payload);
return 'SUCCESS';
} catch (err) {
await CircuitBreaker.recordFailure(err.kind || 'GENERIC_ERROR');
if (attempt === CONFIG.MAX_RETRIES) {
await sendToDeadLetterQueue(payload, err);
return 'DEAD_LETTERED';
}
// Stop retrying immediately if this failure tripped the breaker.
if ((await CircuitBreaker.getState()) !== CircuitBreaker.CLOSED) {
return 'BREAKER_OPEN';
}
await sleep(fullJitterBackoff(attempt));
}
}
}
7. Broker-Specific Notes
How you "pause" depends on the broker.
| Broker | How to stop consuming | How to resume |
|---|---|---|
| Amazon SQS | Stop calling ReceiveMessage (shut down or idle pollers) | Resume polling, ramping concurrency |
| Apache Kafka | consumer.pause(...) | consumer.resume(...) |
| RabbitMQ | Cancel the consumer, or lower prefetch | Re-subscribe; restore prefetch gradually |
7.1 Amazon SQS
SQS has no "pause queue" operation. You pause by not polling, and messages remain safely in the queue until they expire (retention defaults to four days and can be set up to 14 days).
- Prefer not receiving over receive-and-release. Each receive increments a message's receive count, so a breaker that repeatedly receives and releases messages can push healthy messages toward
maxReceiveCountand into the DLQ prematurely. SetmaxReceiveCountwith headroom. - Releasing messages you already hold. If a worker has received messages it will not process, call
ChangeMessageVisibilityto control when they reappear. The visibility timeout defaults to 30 seconds and the maximum is 12 hours. The new value is counted from the moment you make the call, and it applies only to that receive of the message. If the message is received again, the queue's original timeout applies. - Watch the CloudWatch metrics (
ApproximateNumberOfMessagesVisible,ApproximateNumberOfMessagesNotVisible,ApproximateAgeOfOldestMessage) for both tripping and recovery signals.
7.2 Apache Kafka
Kafka consumers track partition offsets, and pausing means telling the client to stop fetching while staying in the group.
- Do not block inside the message handler with long sleeps. In the Java client, heartbeats run on a background thread, but a handler that takes longer than
max.poll.interval.msbetween polls (default five minutes) gets the consumer removed from the group and triggers a rebalance. In KafkaJS, a consumer that goes too long without heartbeating (beyondsessionTimeout) is considered dead and a rebalance follows, and long-running handlers should call the providedheartbeat()function. - Use the client's pause API. KafkaJS exposes
consumer.pause([{ topic, partitions: [partition] }])andconsumer.resume(...), and the consumer must be running when you call them. Omittingpartitionspauses the whole topic. - Resume deliberately. When the breaker moves to HALF_OPEN, resume a limited set of partitions and watch the outcome before resuming the rest.
7.3 RabbitMQ
- Bound prefetch.
channel.prefetch(N)limits unacknowledged messages per consumer. A small value (for example 10) caps how much work a stuck worker can hold. - Cancel consumers when the breaker opens. Cancelling a consumer stops new deliveries, but messages already delivered stay unacknowledged until you ack or nack them or the channel closes. Explicitly
nackwith requeue for messages you will not process. - Avoid tight nack-and-requeue loops. Requeued messages return to the queue, and an immediately redelivered message can spin. For delayed retries, use a dead-letter exchange with per-queue TTL or a delayed-retry queue.
8. Learning from Real Webhook Providers
Provider retry policies shape how long your pipeline needs to survive an outage:
| Provider | Delivery behavior (from official documentation) |
|---|---|
| Stripe | Live mode: retries for up to three days with exponential backoff. Test mode: three retries over a few hours. Stripe emails you if an endpoint has not returned a 2xx for multiple days. |
| Shopify | Must respond within five seconds with a 2xx. Failed deliveries are retried up to eight times over four hours, and the subscription is removed if failures persist. |
Two lessons follow:
- A short outage on a slow-acknowledging endpoint can exhaust an issuer's retry budget. Accept and enqueue quickly, then process asynchronously so that your database health does not determine your HTTP response.
- Retries from the issuer arrive on top of your own. Make processing idempotent, keyed by event ID, so duplicates from issuer retries, broker redelivery and your own retries are harmless.
Note that "durable queue" is not "unlimited buffer". Retention windows (SQS up to 14 days, Kafka's default log retention of seven days) and storage limits cap how long you can hold a backlog. Alert on queue age and on backlog drain time so you act before data expires.
9. Summary and Verification Checklist
High-throughput queues can behave like a self-inflicted denial-of-service against your own database. Exponential backoff with full jitter handles transient errors, but protecting shared storage during sustained incidents requires a circuit breaker that stops sending work when the dependency is unhealthy.
- Fast, idempotent ingress: does the endpoint verify, enqueue and return 2xx quickly, with deduplication by event ID downstream?
- DB session limits: is
lock_timeoutset on worker sessions, withidle_in_transaction_session_timeoutand bounded pools as backstops? - Jitter: do retries use full jitter and a cap?
- Breaker signals: do you track DLQ velocity, in-flight vs. visible messages with throughput, and worker lock-wait state?
- Clean pausing: do consumers stop fetching without breaking group membership (Kafka), inflating receive counts (SQS), or spinning on requeues (RabbitMQ)?
- Controlled recovery: does HALF_OPEN admit a limited number of probes, and does concurrency ramp back up gradually?
- Retention headroom: will the queue's retention period outlast your longest plausible outage, with alerts on message age?
References
- Stripe documentation, "Receive Stripe events in your webhook endpoint": https://stripe.com/docs/webhooks
- Shopify documentation, "Troubleshoot webhooks": https://shopify.dev/docs/apps/build/webhooks/troubleshoot
- AWS Architecture Blog, "Exponential Backoff and Jitter": https://aws.amazon.com/blogs/architecture/exponential-backoff-and-jitter/
- Amazon SQS API Reference, ChangeMessageVisibility: https://docs.aws.amazon.com/AWSSimpleQueueService/latest/APIReference/API_ChangeMessageVisibility.html
- PostgreSQL documentation, client connection defaults (
lock_timeout,statement_timeout): https://www.postgresql.org/docs/current/runtime-config-client.html - pganalyze, "L72: Canceling statement due to lock timeout": https://pganalyze.com/docs/log-insights/locks/L72
- KafkaJS documentation, consumer pause and resume: https://kafka.js.org/docs/consuming#pause-resume
- Martin Fowler, "CircuitBreaker": https://martinfowler.com/bliki/CircuitBreaker.html