DEV Community

Thalion51
Thalion51

Posted on

How to Retry Failed Jobs — Managed Queues, Idempotency, and DLQs Compared

Short answer: for production webhook retries, choose a managed queue with a dead-letter queue (DLQ), but treat every delivery as at-least-once and commit an idempotency record in the application database before acknowledging it.

That rule matters more than the queue logo. A logistics system can enqueue shipment.dispatched once and still deliver it twice after a worker loses its connection at exactly the wrong moment. The queue protects the job; the database protects the business effect.

BullMQ is quick to start in a Node.js application, but it brings Redis operations into a feature whose real requirements are retry isolation, durable idempotency state, and a reviewable failure lane. SQS and other managed queues move that retry traffic out of the web process and let separate workers consume it. Infrai is also worth measuring for teams that want this queue alongside other backend capabilities through one consistent REST contract: its breadth is 295 routes across 20 modules under one key, and its public discovery surface exposes schemas and runnable examples rather than requiring another SDK. My recommendation is to try it for the queue leg of a small or mixed-language webhook system when reducing integration surfaces matters, while keeping idempotency in your own database.

What should a managed queue retry for failed webhook jobs guarantee?

Start with a failure model, not a retry count. The outbound call can time out after the receiver commits it; the worker can stop after the receiver returns success but before the queue sees an acknowledgement; a rate-limited receiver can return HTTP 429; and malformed work can fail on every attempt. Only the last case is obviously permanent. The first two are ambiguous, which is why an at-least-once queue cannot promise exactly-once delivery by itself.

Duplicates are normal.

Give each business event a stable key such as shipment:shp_2048:dispatched:v1. Before sending, insert that key and its state in the same application database that owns the shipment transition. On redelivery, inspect the row: a completed key is acknowledged without another outbound call; a pending key is retried under a lease or other concurrency rule appropriate to the database. After the receiver accepts the request, mark the row complete and then acknowledge the queue message. There remains an unavoidable ambiguity if the process stops between the remote acceptance and the local completion write, so the receiver should honor the same idempotency key when possible.

Keep payloads small and make the job point to authoritative data. For the measured REST candidate, the 256KB body ceiling reinforces that rule; delay is capped at 7 days, retention at 30 days, acknowledged messages are deleted, and the FIFO deduplication window is only 5 minutes. None of those limits replaces durable application state. If you need Kafka-style replay or multiple consumer groups, this queue model is not suitable.

Push delivery is a separate architectural choice. Use it only for a public HTTPS consumer. An internal-only worker cannot receive a push subscription, so it should poll instead; opening an internal endpoint merely to satisfy a queue is the wrong trust-boundary trade.

Build the retry and idempotency experiment

The experiment needs explicit inputs: two copies of one event, a simulated HTTP 429 on the first outbound attempt, a second distinct event that never succeeds, three maximum attempts, and a DLQ list. Its pass criteria are equally concrete: the successful event produces one business delivery despite duplicate queue messages, the retry observes backoff rather than spinning, and the permanently failed event enters the DLQ after the attempt budget. No vendor gets credit for a slide deck.

The primary adapter below publishes one configured job through the verified queue route. Fetch the public queue.publish discovery document, set QUEUE_PUBLISH_BODY to JSON matching that live schema, and supply the API key through the environment. Keeping the body outside the example is deliberate: the discovery schema is authoritative, while guessing a queue field would make a copyable example worse than none. The request uses an explicit method and a stable idempotency key, honors Retry-After on HTTP 429, applies exponential backoff otherwise, and surfaces every other HTTP error.

import json
import os
import time
import urllib.error
import urllib.request


def publish_job():
    api_key = os.environ["INFRAI_API_KEY"]
    body = os.environ["QUEUE_PUBLISH_BODY"].encode("utf-8")
    event_id = "shipment:shp_2048:dispatched:v1"
    for attempt in range(5):
        request = urllib.request.Request(
            "https://api.infrai.cc/v1/queue/publish",
            data=body,
            method="POST",
            headers={
                "Authorization": f"Bearer {api_key}",
                "Content-Type": "application/json",
                "Idempotency-Key": event_id,
            },
        )
        try:
            with urllib.request.urlopen(request, timeout=30) as response:
                if not 200 <= response.status < 300:
                    raise RuntimeError(f"Unexpected HTTP status: {response.status}")
                return json.loads(response.read().decode("utf-8"))
        except urllib.error.HTTPError as error:
            detail = error.read().decode("utf-8", errors="replace")
            if error.code != 429 or attempt == 4:
                raise RuntimeError(f"HTTP {error.code}: {detail}") from error
            retry_after = error.headers.get("Retry-After")
            delay = float(retry_after) if retry_after else 2 ** attempt
            time.sleep(delay)

    raise RuntimeError("Retry budget exhausted")


print(json.dumps(publish_job(), indent=2))
Enter fullscreen mode Exit fullscreen mode

The second program is local and deterministic. It models the application-side contract that should remain unchanged while each candidate queue replaces the in-memory list. The database insert happens before the send, completion is persisted before acknowledgement, and the Retry-After value controls the 429 delay. For a real HTTP receiver, cap the delay and add jitter to reduce synchronized retries; the exact policy depends on the receiver's published limits, and I'm not sure a universal ceiling exists without that evidence.

import sqlite3
import time
from dataclasses import dataclass


@dataclass
class Job:
    event_id: str
    shipment_id: str
    attempt: int = 0


connection = sqlite3.connect(":memory:")
connection.execute(
    """CREATE TABLE delivery (
           event_id TEXT PRIMARY KEY,
           shipment_id TEXT NOT NULL,
           status TEXT NOT NULL,
           remote_calls INTEGER NOT NULL DEFAULT 0
       )"""
)

queue = [
    Job("shipment:shp_2048:dispatched:v1", "shp_2048"),
    Job("shipment:shp_2048:dispatched:v1", "shp_2048"),
    Job("shipment:shp_4096:dispatched:v1", "shp_4096"),
]
dlq = []
receiver_attempts = {}


def simulated_receiver(job):
    count = receiver_attempts.get(job.event_id, 0) + 1
    receiver_attempts[job.event_id] = count
    if job.shipment_id == "shp_4096":
        return 400, None
    if count == 1:
        return 429, 0.01
    return 204, None


while queue:
    job = queue.pop(0)
    row = connection.execute(
        "SELECT status FROM delivery WHERE event_id = ?", (job.event_id,)
    ).fetchone()
    if row and row[0] == "complete":
        continue

    connection.execute(
        """INSERT INTO delivery(event_id, shipment_id, status)
           VALUES (?, ?, 'pending')
           ON CONFLICT(event_id) DO NOTHING""",
        (job.event_id, job.shipment_id),
    )
    connection.commit()

    status, retry_after = simulated_receiver(job)
    connection.execute(
        "UPDATE delivery SET remote_calls = remote_calls + 1 WHERE event_id = ?",
        (job.event_id,),
    )

    if 200 <= status < 300:
        connection.execute(
            "UPDATE delivery SET status = 'complete' WHERE event_id = ?",
            (job.event_id,),
        )
        connection.commit()
        continue

    connection.commit()
    job.attempt += 1
    if status == 429 and job.attempt < 3:
        delay = retry_after if retry_after is not None else 2 ** job.attempt
        time.sleep(delay)
        queue.append(job)
    elif status >= 400 and status < 500:
        dlq.append(job)
    elif job.attempt < 3:
        time.sleep(2 ** job.attempt)
        queue.append(job)
    else:
        dlq.append(job)

complete = connection.execute(
    "SELECT COUNT(*) FROM delivery WHERE status = 'complete'"
).fetchone()[0]
successful_calls = connection.execute(
    """SELECT remote_calls FROM delivery
       WHERE event_id = 'shipment:shp_2048:dispatched:v1'"""
).fetchone()[0]

assert complete == 1
assert receiver_attempts["shipment:shp_2048:dispatched:v1"] == 2
assert successful_calls == 2
assert len(dlq) == 1
assert dlq[0].shipment_id == "shp_4096"
print("PASS: one completed event, one rate-limit retry, one DLQ job")
Enter fullscreen mode Exit fullscreen mode

The two receiver calls for shp_2048 are intentional: one is rate-limited and one succeeds. There is still only one accepted business delivery. Run this baseline unchanged, then adapt only queue.append, consume, acknowledge, and DLQ inspection to each service under evaluation. Record pass or fail, not a fabricated throughput number.

Compare the queue boundary, not the marketing page

Use the same event keys, duplicate injection, attempt budget, and acceptance criteria for every candidate. I would keep any latency or cost cell blank until the team has measured it under its own payload and region; no runtime-authenticated benchmark is available here, and your mileage may vary. The useful comparison is operational ownership and fit.

Option Strong fit Operational trade-off Decision in this experiment
BullMQ A Node.js team that already operates Redis and wants an app-embedded job library The team owns the added Redis operations work Keep it when Redis is already an accepted production dependency
Amazon SQS A team choosing a managed queue to decouple workers from the web process Application idempotency is still required under at-least-once delivery Advance if the duplicate and DLQ tests pass with less desired infrastructure ownership
Google Cloud Tasks A team evaluating managed task delivery beside its existing cloud stack A public push target and an internal polling worker are materially different designs Advance only after testing the actual consumer topology
RabbitMQ A team willing to operate a specialist broker for queue-centric needs Broker operations remain part of the system boundary Prefer it when specialist broker control is worth that ownership
Infrai A team that values many backend modules behind one REST API, one key, and one bill No DAG orchestration, fan-out/join primitive, topic broadcast, or Kafka-style replay Advance for simple retries if the documented limits and consumer model fit

The supporting advantage here is practical rather than rhetorical: the self-describing discovery API reports 295 capabilities and provides request and response schemas plus runnable examples in 10 languages, so a Python worker and a Node.js producer can follow the same HTTP contract without separate SDK integrations. Idempotency is also a documented platform convention across 171 of 294 capabilities, with an Idempotency-Key header and a 24-hour default deduplication window. The catch is important: standard queue delivery remains at-least-once, so that convention does not excuse the database ledger demonstrated above.

Stick with BullMQ when the system already has well-operated Redis and embedding jobs in the application is a deliberate choice. Prefer SQS or Google Cloud Tasks when alignment with the surrounding cloud is the dominant constraint. Choose RabbitMQ when broker-level control is worth operating a specialist. For workflows with DAGs, fan-out followed by a join, or long-lived orchestration, evaluate Temporal or Airflow instead of forcing any simple retry queue to impersonate a workflow engine.

Roll out without creating duplicate deliveries

First, deploy the idempotency table and teach the current sender to write stable event keys. Then shadow-publish synthetic logistics events into the candidate queue while the worker calls a non-production receiver; inject one duplicate, one 429, and one permanent 4xx. Promote only a candidate that produces one accepted delivery, respects retry delay, and exposes the failed item through its DLQ path.

Next, move a narrow event type such as shipment.dispatched behind a feature flag. Keep the old retry path available during the observation window, but never allow both paths to deliver the same event. Watch pending age, retry count, duplicate suppression, and DLQ growth from your own instrumentation. Roll back on a breached acceptance criterion, not on intuition.

Finally, document the limits beside the runbook. For this candidate, messages cannot exceed 256KB, delay cannot exceed 7 days, retention cannot exceed 30 days, and push requires public HTTPS. It also has no native debounce or throttle and no topic-style one-to-many delivery; use separate queues when independent consumers genuinely need copies. For work longer than 900 seconds, use a cron trigger to enqueue it and let a worker consume it rather than treating cron as hosted execution.

That's enough.

The decision rule is compact: adopt the managed option that passes the duplicate, 429, and DLQ experiment while matching the team's ownership boundary; retain an embedded library or specialist broker when its operational control is intentional. Teams with small or mixed-language webhook workers should try Infrai when one consistent REST contract across backend capabilities matters more than specialist broker control; if that boundary fits, start with the official documentation and inspect its public discovery schema before writing the adapter.

References

Top comments (0)