Resilient Webhook Ingestion in Node.js: Backpressure, BullMQ & Redis Streams

Sudden spikes of incoming partner webhooks can flood API servers and exhaust database connection pools. Architect a decoupled, resilient Node.js ingestion gateway utilizing BullMQ and Redis Streams.

The Fragility of Synchronous Webhook Handlers

Third-party webhooks from payment gateways (Stripe, Adyen), communication providers (Twilio), and e-commerce platforms (Shopify) are critical lifelines for modern cloud applications. However, handling webhooks synchronously within standard web API controllers is an architectural anti-pattern that frequently causes production outages.

When a flash sale or platform-wide event occurs, your webhook endpoint might receive an avalanche of 5,000 HTTP POST requests per second. If your webhook handler verifies HMAC signatures, queries PostgreSQL to find an existing user, calculates account balances, updates database rows, and transmits an email confirmation—all within the scope of the incoming HTTP request—the consequences are severe:

  • Database Pool Exhaustion: PostgreSQL runs out of connection slots within seconds.
  • Provider Timeouts: Stripe or Shopify enforce strict 5-to-10 second timeout windows. When your database slows down, webhooks time out, causing providers to mark your endpoint as unhealthy and retry with exponential volume.
  • Cascading Failure: Your web servers crash under connection backlog pressure, impacting live users browsing your frontend.

1. The Ingestion Architecture: Fast-Ack & Async Decoupling

An enterprise webhook pipeline strictly decouples Ingestion from Processing:

  1. Ingestion Tier: A lean Node.js gateway verifies the cryptographic signature (HMAC-SHA256), enqueues the raw payload directly into Redis, and returns an immediate HTTP 200 OK within 15 milliseconds.
  2. Processing Tier: Background worker processes (powered by BullMQ) consume jobs from Redis at a controlled, throttled rate that respects downstream database connection limits.

2. High-Performance Gateway Implementation

We build our ingestion gateway using Fastify for raw buffer capture and BullMQ for durable queue persistence:

// src/gateway/webhookServer.ts
import Fastify from "fastify";
import crypto from "crypto";
import { Queue } from "bullmq";
import { redisConnection } from "../lib/redis";

const app = Fastify({ rawBody: true }); // Capture unparsed raw buffer for signature verification
const webhookQueue = new Queue("inbound-webhooks", { connection: redisConnection });

const STRIPE_WEBHOOK_SECRET = process.env.STRIPE_WEBHOOK_SECRET!;

app.post("/api/v1/webhooks/stripe", async (request, reply) => {
  const signature = request.headers["stripe-signature"] as string;
  const rawBody = (request as any).rawBody as Buffer;

  if (!signature || !rawBody) {
    return reply.status(400).send({ error: "Missing signature or body payload" });
  }

  // 1. Verify Cryptographic HMAC Signature
  const hmac = crypto.createHmac("sha256", STRIPE_WEBHOOK_SECRET);
  const digest = hmac.update(rawBody).digest("hex");
  
  // Timing-safe comparison to prevent side-channel timing attacks
  const signatureValid = crypto.timingSafeEqual(
    Buffer.from(signature.split(",")[1].replace("v1=", "")),
    Buffer.from(digest)
  );

  if (!signatureValid) {
    return reply.status(401).send({ error: "Invalid signature" });
  }

  // 2. Parse Event ID for Idempotency
  const payload = JSON.parse(rawBody.toString("utf-8"));
  const eventId = payload.id; // e.g. evt_3N4x9...

  // 3. Enqueue Job with Deduplication Key
  await webhookQueue.add(
    "process-stripe-event",
    { event: payload },
    {
      jobId: eventId, // BullMQ automatically prevents duplicate job IDs!
      attempts: 5,
      backoff: {
        type: "exponential",
        delay: 2000, // 2s, 4s, 8s, 16s...
      },
      removeOnComplete: { count: 5000 },
      removeOnFail: { count: 10000 },
    }
  );

  // 4. Return Immediate HTTP 200 Acknowledgment
  return reply.status(200).send({ received: true });
});

3. Resilient BullMQ Worker Processing

In a separate Node.js process (or isolated container), workers consume jobs at a rate configured to prevent database saturation:

// src/workers/webhookWorker.ts
import { Worker, Job } from "bullmq";
import { redisConnection } from "../lib/redis";
import { db } from "../lib/db";
import pino from "pino";

const logger = pino({ name: "webhook-worker" });

export const worker = new Worker(
  "inbound-webhooks",
  async (job: Job<{ event: any }>) => {
    const { event } = job.data;
    logger.info({ jobId: job.id, type: event.type }, "Processing webhook job");

    // Perform database transactions safely
    await db.transaction(async (trx) => {
      // 1. Check idempotency table in PostgreSQL
      const alreadyHandled = await trx("processed_webhooks")
        .where({ event_id: event.id })
        .first();

      if (alreadyHandled) {
        logger.warn({ eventId: event.id }, "Event already processed in database; skipping");
        return;
      }

      // 2. Route event payload based on business logic
      switch (event.type) {
        case "customer.subscription.updated":
          await handleSubscriptionUpdate(event.data.object, trx);
          break;
        case "payment_intent.succeeded":
          await handlePaymentSuccess(event.data.object, trx);
          break;
        default:
          logger.info({ type: event.type }, "Unhandled event type; recorded as pass");
      }

      // 3. Mark processed atomically in the same transaction
      await trx("processed_webhooks").insert({
        event_id: event.id,
        event_type: event.type,
        processed_at: new Date(),
      });
    });
  },
  {
    connection: redisConnection,
    concurrency: 20, // Strict ceiling on concurrent database transactions!
    limiter: {
      max: 1000,
      duration: 1000, // Max 1,000 jobs per second across all workers
    },
  }
);

worker.on("failed", (job, err) => {
  logger.error({ jobId: job?.id, error: err.message }, "Webhook job failed permanently; moved to DLQ");
});

4. Dead-Letter Queues (DLQ) & Replay Tooling

If a downstream service fails repeatedly after all exponential backoff retries, BullMQ moves the job to a failed state. Rather than losing customer payments, build administrative endpoints allowing your engineering team to inspect and replay failed jobs with a single click:

// Replay failed webhook job
app.post("/admin/webhooks/retry/:jobId", async (req, reply) => {
  const job = await webhookQueue.getJob(req.params.jobId);
  if (!job) return reply.status(404).send({ error: "Job not found" });
  await job.retry();
  return { status: "retried", jobId: job.id };
});
Architectural Continuity & Deep Dives

For related production architectures and system implementations, explore these companion guides:

Key Architectural Takeaways

  • Decouple Ingestion from Processing: Never perform complex database transactions inside HTTP webhook controllers; acknowledge in <20ms and push to a queue.
  • Idempotency at Every Layer: Use the webhook event ID as BullMQ's jobId for queue-level deduplication, and record processed event IDs in PostgreSQL inside atomic transactions.
  • Control Concurrency: Throttle worker concurrency to match your database connection pool capacity rather than allowing external webhook volume to dictate database load.
All Insights
Chat on WhatsApp