InstaWebhook
October 5, 2026By InstaWebhook TeamIncident Response

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

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.

Code example
+----------------+   +--------------------+   +---------------+   +--------------+
| 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:

  1. Push the message back onto the main queue or a retry topic.
  2. Retry at fixed intervals, or with uncapped exponential backoff and no jitter.
  3. 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:

  1. Transaction A updates orders where id = 1042, taking a row-level lock that is held until commit or rollback.
  2. A storage slowdown stretches Transaction A from about 10 ms to several seconds.
  3. 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 Lock wait event in pg_stat_activity.
  4. Blocked workers hold their connections, so the pool (HikariCP, pg.Pool, or a PgBouncer limit) is exhausted.
  5. Webhooks for unrelated orders (1043, 1044) cannot get a connection, fail, and re-enter the retry flow.
  6. The retry flow multiplies lock requests and connection attempts until the database saturates.
Code example
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.

DimensionExponential backoff with jitterCircuit breaker
ScopeIndividual message or callThe whole path to a failing resource
MechanismDelays re-execution of one payloadStops or throttles all calls to the resource
Best forTransient network glitches, brief racesSustained outages, lock exhaustion, overload
Load on the dependencyStill sends work, just spread outSheds 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.

Code example
                         +------------------------+
                         | 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:

  1. CLOSED (normal): workers pull at full capacity, and failure signals stay below thresholds.
  2. 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.
  3. 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​+Mvisible​Minflight​​

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:

Code example
{
  "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 SQLSTATE 55P03 (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 a statement_timeout shorter than lock_timeout makes 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 NX lock 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.
Code example
/**
 * 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.

BrokerHow to stop consumingHow to resume
Amazon SQSStop calling ReceiveMessage (shut down or idle pollers)Resume polling, ramping concurrency
Apache Kafkaconsumer.pause(...)consumer.resume(...)
RabbitMQCancel the consumer, or lower prefetchRe-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 maxReceiveCount and into the DLQ prematurely. Set maxReceiveCount with headroom.
  • Releasing messages you already hold. If a worker has received messages it will not process, call ChangeMessageVisibility to 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.

  1. 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.ms between 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 (beyond sessionTimeout) is considered dead and a rebalance follows, and long-running handlers should call the provided heartbeat() function.
  2. Use the client's pause API. KafkaJS exposes consumer.pause([{ topic, partitions: [partition] }]) and consumer.resume(...), and the consumer must be running when you call them. Omitting partitions pauses the whole topic.
  3. 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

  1. 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.
  2. 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 nack with requeue for messages you will not process.
  3. 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:

ProviderDelivery behavior (from official documentation)
StripeLive 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.
ShopifyMust 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_timeout set on worker sessions, with idle_in_transaction_session_timeout and 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