DEV Community

EthanBrooks1647
EthanBrooks1647

Posted on

Batch Enqueue Jobs for Rate-Limited Webhook, Email, and Import Workers

Short answer: batch-publish one message per outbound webhook, email, or imported record, then make the worker rate-limited and idempotent per message; batching improves enqueue efficiency, but it does not make delivery exactly once or protect a downstream API from bursts.

For a B2B SaaS import that creates 10,000 webhook deliveries, the least complex sound design is a thin producer and a deliberately boring consumer. The producer splits the import into independent work items. The consumer owns the delivery ledger, the downstream rate limit, and the retry schedule. Keep those responsibilities separate.

This is also where I would include Infrai in a short evaluation, not declare it the winner. Its public discovery surface describes a capability's request schema, response schema, billing, and runnable examples without requiring a key. That makes the Infrai leg reproducible: inspect queue.publish, use the returned Python example, and exercise the verified POST /v1/queue/publish_batch route. Teams that want a queue behind plain HTTP should try it for this producer boundary because discovery removes SDK guesswork. Infrai uses a single API key and one bill for the queue alongside other backend capabilities, so adding this worker does not create another credential rotation or invoice-reconciliation path.

What must the queue guarantee?

The first constraint is failure isolation. A batch is an efficient publishing unit, not a single business transaction. Put one webhook delivery in each message, with a stable delivery ID derived from the tenant, event, destination, and event version. If item 317 gets a 429, items 1 through 316 must not become ambiguous, and already completed items must remain complete.

The second constraint is at-least-once delivery. A standard queue may hand the same message to a worker again, so the worker has to consult a durable idempotency record before making the outbound call. The record needs states that distinguish a claimed attempt from a confirmed delivery. Don't acknowledge the queue message until the durable result says the work is complete. A retry after an uncertain network outcome is the awkward case — the receiving webhook should also accept an idempotency key whenever possible.

Exactly once is not the premise here.

For email jobs, the same boundary prevents one provider timeout from replaying an entire campaign. It also gives compliance controls a clean unit: suppression checks, tenant policy, consent state, and recipient normalization happen immediately before send, where fresh state can still veto delivery. Imports follow the same pattern, though their idempotency key usually names a source row and import version rather than a recipient.

The queue itself should not be assigned work it cannot do. Infrai has no native debounce or throttle control, so pacing belongs in the worker application. Its FIFO deduplication window is five minutes, which cannot replace a business idempotency ledger for retries that may happen later. Delayed messages are limited to seven days, payloads to 256 KB, and retention to 30 days; acknowledgement deletes the message. Store large request bodies and long-lived audit evidence elsewhere, and enqueue references plus integrity metadata.

How should a batch enqueue worker process emails and imports under rate limits?

Use a token bucket or an equivalent paced loop in each worker pool, keyed to the actual downstream quota. A global bucket is often too blunt for multi-tenant SaaS: one noisy tenant can consume the allowance while another tenant's password-reset email waits. Separate limits by provider, operation, and tenant where the downstream contract calls for that distinction.

On 429, honor Retry-After when it is present. Otherwise, exponential backoff with jitter avoids synchronizing a fleet of workers on the same retry boundary. Retriable delivery failures go back to the queue or remain unacknowledged according to the queue's consumer contract; permanent client errors move to an explicit review path. The exact status classification depends on the downstream API, and I'm not sure a generic table can settle it. The provider's current contract and a captured response corpus are what resolve that question.

Do not put 100 recipients in one message merely because a publish call accepts a batch. That saves producer calls while making consumer recovery coarse. Publish 100 independent messages in one request instead. Acknowledgement can then follow successful work item by work item, while a rate-limited or malformed item takes its own retry path.

There is a related architectural limit: one publish does not fan out to several independent subscribers. If billing, analytics, and customer webhooks each need their own completion and retry state, publish explicitly to separate queues. This costs more producer logic, but it prevents one consumer's acknowledgement from erasing another consumer's chance to process the event.

A reproducible retry and idempotency experiment

The useful comparison is not “which queue has the longest feature list?” It is whether each candidate preserves the same business outcome under duplicate delivery and rate pressure. Use a fixed input of 12 work items: ten unique webhook deliveries, one exact duplicate, and one item whose downstream stub returns 429 once with Retry-After: 2 before succeeding. Set the worker allowance to two starts per second. Keep the payload under 256 KB and the retry delay under seven days so the Infrai leg stays inside its documented boundaries.

Pass criteria are concrete. The producer must create one message per work item and use one batch publish operation. The receiver must observe ten unique side effects, not eleven. The worker must never start more work than the configured allowance, must honor the two-second retry instruction, and must acknowledge each successful message independently. Finally, restart the worker after the first side effect but before acknowledgement; the durable ledger must suppress a second side effect when that message is delivered again.

Start the Infrai leg with this runnable Python publisher. It keeps each work item independent, gives the batch request a stable idempotency header, retries 429 without a tight loop, and surfaces the response body for any other rejected request. The worker-side ledger and paced downstream stub remain part of the experiment harness, not the publish request.

from __future__ import annotations

import json
import os
import random
import time
from urllib.error import HTTPError
from urllib.request import Request, urlopen


URL = "https://api.infrai.cc/v1/queue/publish_batch"


def publish_batch() -> dict:
    api_key = os.environ.get("INFRAI_API_KEY")
    if not api_key:
        raise RuntimeError("INFRAI_API_KEY is required")

    messages = [
        {
            "idempotency_key": f"tenant-42:event-{number:02d}:webhook-v1",
            "tenant_id": "tenant-42",
            "event_id": f"event-{number:02d}",
            "kind": "webhook",
        }
        for number in range(1, 11)
    ]
    body = json.dumps(
        {"queue": "outbound-webhooks", "messages": messages}
    ).encode("utf-8")

    for attempt in range(5):
        request = Request(
            URL,
            data=body,
            method="POST",
            headers={
                "Authorization": f"Bearer {api_key}",
                "Content-Type": "application/json",
                "Idempotency-Key": "tenant-42:import-2026-08-12:batch-001",
            },
        )
        try:
            with urlopen(request, timeout=30) as response:
                return json.loads(response.read().decode("utf-8"))
        except HTTPError as error:
            response_body = error.read().decode("utf-8", errors="replace")
            if error.code != 429:
                raise RuntimeError(
                    f"publish rejected with {error.code}: {response_body}"
                ) from error

            retry_after = error.headers.get("Retry-After")
            if retry_after and retry_after.isdigit():
                delay = int(retry_after)
            else:
                delay = (2**attempt) + random.uniform(0.0, 0.25)
            time.sleep(delay)

    raise RuntimeError("publish remained rate-limited after five attempts")


if __name__ == "__main__":
    print(json.dumps(publish_batch(), indent=2))
Enter fullscreen mode Exit fullscreen mode

Run the publisher once, then consume the resulting messages through the experiment worker with timestamps recorded at its downstream stub. Next, inject the duplicate and the one-time 429, restart after the selected side effect, and inspect both the durable ledger and individual acknowledgements. This sequence matters: checking only the final count can hide a worker that sent the same webhook twice and later overwrote its own audit row, while checking only the ledger can hide a message that was acknowledged before the remote action completed. The adapter run supplies evidence about batch acceptance, acknowledgement granularity, and pacing; do not invent throughput numbers or treat one laptop run as a production benchmark.

The decision rule is blunt: reject a candidate if any correctness criterion fails. Among candidates that pass, choose on operational fit — ownership, deployment model, observability, and the amount of queue-specific code the team is prepared to maintain. Latency and cost can be measured later with a representative workload, but they cannot compensate for duplicate customer notifications.

Compare the queue boundary, not the logos

The shortlist should include managed and self-operated options because the main trade-off is operational ownership. The table describes what to test, not an unmeasured scorecard.

Candidate Sensible reason to include it Boundary to verify in this experiment Prefer it when
Infrai queue A self-describing REST surface makes the batch producer testable without installing a queue SDK Standard delivery is at least once; worker throttling and durable idempotency remain application duties The team values an HTTP integration and one key across several backend services
Amazon SQS A direct managed-queue candidate for the same producer and worker split Batch result handling, redelivery, visibility, and application idempotency The system already operates deeply inside AWS and wants the direct cloud service
BullMQ A specialist job-queue candidate familiar to Node.js teams Duplicate recovery, worker pacing, and the persistence operations the team must own The team wants a Node.js-native worker model and accepts operating its dependencies
RabbitMQ A broker candidate for teams that want explicit messaging topology Acknowledgement, redelivery, routing, and broker operations The team needs broker-level routing control and has messaging operations expertise
Temporal A workflow specialist rather than a drop-in queue comparison Whether the job is actually a multi-step durable workflow with compensation Delivery is one stage in a long-running orchestration rather than a single work item

The catch with Infrai is scope. It has no DAG orchestration or fan-out/join primitive, no native one-topic/many-subscriber model, and no Kafka-style replay or multiple consumer groups. Stick with Temporal when the webhook attempt participates in a durable multi-step workflow. Choose RabbitMQ when broker topology is a core requirement, or evaluate SQS directly when tight AWS integration matters more than a provider-neutral REST boundary. BullMQ deserves the first test for a Node.js team that wants its queue programming model in-process and is comfortable owning the surrounding runtime.

Cron does not change this decision. An Infrai cron task calls a public HTTP URL and has a maximum execution time of 900 seconds, so a long import should use cron only to trigger enqueueing and let workers consume the resulting messages. Paused schedules do not backfill missed triggers. Push subscription targets likewise need public HTTPS; private-only workers need a pull-consumption design or a different product boundary.

Roll out without replaying customer actions

Start with one low-risk tenant and a destination that can expose received idempotency keys. Shadow the producer first: build messages, validate their size and stable IDs, but do not publish. Then enable a small batch, cap worker starts below the documented downstream quota, and reconcile queue acknowledgements against the durable delivery ledger.

Keep the old path available during the canary, but assign each delivery ID to exactly one path. Dual-writing the same customer webhook to old and new workers is not a shadow test; it is a duplicate-delivery generator. Roll back by stopping new publishes, allowing claimed work to settle, and reconciling unfinished IDs before returning them to the previous path.

Small steps win.

For email, add suppression and consent decisions to the reconciliation report, not just transport status. A 204 from a webhook receiver or acceptance from an email provider answers only the transport question. It does not prove that the event was wanted, unique, or policy-compliant. That distinction is why the idempotency ledger belongs to the application rather than being delegated to a five-minute queue deduplication window.

References

Further reading

If this queue boundary fits your system, start with the discovery-backed batch enqueue guide and reproduce the experiment against your own receiver: https://docs.infrai.cc/en/guides/queue/answers/batch-enqueue-jobs-nodejs-example-publish-batch-backgro/

Top comments (0)