Deduplicating Jobs in BullMQ

Many duplicate jobs are created before a worker ever sees them: a user double-clicks, a webhook is delivered twice, an event fires on every keystroke. BullMQ offers several enqueue-time deduplication mechanisms, and choosing the right one depends on whether you want the first request, the last one, or at most one per time window. This guide covers them as part of BullMQ for Node.js Ecosystems in Backend Frameworks & Worker Scaling.

Problem Statement

A collaboration app enqueues a "reindex document" job every time a document is saved. Autosave fires every few seconds while someone types, so a busy document generates hundreds of reindex jobs per hour, each re-reading and re-indexing the whole document; search indexing falls hours behind. Separately, a billing webhook from a payment provider arrives twice for some events, and the "record payment" job runs twice. You want at most one pending reindex per document that always uses the latest content, exactly one payment job per provider event, and a clear rule for when each mechanism applies.

Prerequisites

  • BullMQ 5.x (the deduplication option with ttl, extend, and replace is available in 5.x releases; check your minor version's documentation).
  • A natural key for each kind of duplicate: document id for reindexing, provider event id for payments.
  • Idempotent processors — enqueue-time deduplication narrows duplicates but cannot remove all of them.

Step 1 — Use a Custom Job Id for Exactly-One Semantics

A custom jobId makes add idempotent for as long as a job with that id exists in the queue. A second add with the same id is ignored and returns the existing job.

// Payment webhooks: one job per provider event, ever (while the job is retained)
await payments.add("record-payment", payload, {
  jobId: `stripe:${event.id}`,                 // provider's event id: stable across redeliveries
  removeOnComplete: { age: 7 * 24 * 3600 },    // keep completed jobs 7 days so late duplicates are ignored
  removeOnFail: { age: 14 * 24 * 3600 },
});

The retention settings matter: once the job is removed, a later add with the same id creates a new job. Keep completed jobs at least as long as duplicates can arrive (providers commonly redeliver for up to a few days). Job ids must not contain : in some older versions and should not be plain integers, which can collide with BullMQ's own counter-generated ids — prefixing with a source avoids both.

jobId dedup lasts as long as the job The first webhook adds job stripe:evt_1, which runs and completes. A redelivered webhook two hours later adds the same id and is ignored because the completed job is still retained. If completed jobs were removed immediately, the redelivery would create a second job and record the payment twice. jobId stripe:evt_1 and a redelivery 2 hours later kept 7 days completed job retained: redelivery ignored removed at once job 1 done redelivery creates job 2: paid twice Retention must outlast the longest redelivery window of the source.

Step 2 — Use Deduplication with a TTL for "At Most One per Window"

The deduplication option keys duplicates by an id you choose, without making that id the job id, and can expire the key after a TTL. While the key is active, further add calls with the same deduplication id are ignored (and emit a deduplicated event).

// "Simple" mode: ignore duplicates while the first job is waiting or active
await reindex.add("reindex", { docId }, { deduplication: { id: `reindex:${docId}` } });

// "Throttle" mode: at most one job per doc per 60 s, regardless of state
await reindex.add("reindex", { docId }, { deduplication: { id: `reindex:${docId}`, ttl: 60_000 } });

Without a TTL, the deduplication key lasts until the job completes or fails; with a TTL, it lasts for that duration. Throttle mode suits work where one run per interval is enough and the first request's data is fine.

Step 3 — Debounce: Keep Only the Latest Pending Job

For autosave-driven reindexing, the first request's data is exactly wrong — you want the latest content. Debounce mode replaces the pending job's data with each new request and optionally pushes its delay back, so the job runs once, after edits pause, with the latest state.

await reindex.add("reindex", { docId, version }, {
  delay: 5_000,                                           // wait 5 s after the last save
  deduplication: {
    id: `reindex:${docId}`,
    ttl: 5_000,
    extend: true,                                         // each new save extends the window
    replace: true,                                        // and replaces the pending job's data
  },
});

While someone types, each save replaces the delayed job and resets the 5-second timer; once they stop, one reindex runs with the latest version. Hundreds of jobs per hour per document become a handful. The processor should still read the current document from the database rather than trusting version blindly, so a job that runs slightly early never indexes stale content.

Debounce collapses a burst into one job A user saves a document every three seconds for thirty seconds. Without deduplication, ten reindex jobs are enqueued and all run. With debounce mode, each save replaces the single pending delayed job and extends its five-second delay, so exactly one reindex runs five seconds after the last save, using the latest version. Ten saves in 30 seconds no dedup 10 reindexes debounce one delayed job, data replaced and timer extended on every save 1 run The single run uses the latest version, five seconds after typing stops.

Step 4 — Pick the Mechanism by What You Want to Keep

Goal Mechanism Keeps Example
Exactly one job per external event jobId + long retention The first Payment webhooks
One job at a time per key deduplication without TTL The first, until it finishes "Sync account"
At most one per time window deduplication with ttl The first in each window Rate-limited refresh
Latest state, after activity settles deduplication with ttl, extend, replace + delay The last Autosave reindex

The key itself deserves as much thought as the mode. It must include every field that makes two requests genuinely different, and nothing that varies between duplicates. reindex:${docId} is right for reindexing because any two reindexes of one document are interchangeable; a key that included a timestamp or request id would never match and would deduplicate nothing. For payment events, the provider's event id is right and the payment amount is wrong — two distinct payments of the same amount must both be recorded. Write the key construction in one helper per job type, next to the retry settings, so it is reviewed with the same care as the processor.

Choosing by "what do I keep" avoids the common mistake of throttling a job whose value is in its latest data, or debouncing a job where every request must be honoured.

Step 5 — Observe Deduplication and Keep Processors Idempotent

Count how often deduplication kicks in; a sudden change indicates a producer bug (more duplicates) or a broken key (fewer).

const events = new QueueEvents("reindex", { connection });
events.on("deduplicated", ({ jobId, deduplicationId }) =>
  dedupedTotal.inc({ queue: "reindex" }));

Enqueue-time deduplication does not cover everything: a job can be redelivered after a worker crash, or run twice after its deduplication key expired. Processors must still be idempotent — for payments, a unique constraint on the provider event id in the payments table, as in idempotent consumers with Postgres unique constraints.

const worker = new Worker("payments", async (job) => {
  await db.query(
    `INSERT INTO payments (provider_event_id, amount_cents, account_id)
     VALUES ($1, $2, $3) ON CONFLICT (provider_event_id) DO NOTHING`,
    [job.data.eventId, job.data.amount, job.data.accountId]);
}, { connection });

Step 6 — Remove a Deduplication Key When Needed

Sometimes a new request must run even though a deduplicated job exists — an admin forces a reindex. Remove the key (or the job) and add again.

async function forceReindex(docId: string) {
  await reindex.removeDeduplicationKey(`reindex:${docId}`);   // name varies by version
  await reindex.add("reindex", { docId }, { deduplication: { id: `reindex:${docId}` } });
}

Provide this as an explicit admin action rather than bypassing deduplication in normal code paths.

Two layers against duplicates Duplicate requests first hit enqueue-time deduplication by job id or deduplication key, which blocks most of them. Duplicates that get through — redeliveries after a worker crash, or requests after a key expired — reach the processor, where a unique constraint or idempotency key makes the second execution a no-op. Belt and braces duplicate requests retries, redeliveries enqueue dedup blocks most idempotent processor absorbs the rest Enqueue dedup saves work; idempotency guarantees correctness. Neither replaces the other.

Verification

it("debounces autosave into one reindex with the latest version", async () => {
  for (let v = 1; v <= 10; v++) {
    await reindex.add("reindex", { docId: "d1", version: v },
      { delay: 200, deduplication: { id: "reindex:d1", ttl: 200, extend: true, replace: true } });
    await sleep(50);
  }
  await sleep(600);
  expect(indexer.runsFor("d1")).toEqual([10]);
});

Gotchas & Edge Cases

Deduplication ids vs job ids. A custom jobId and a deduplication.id are independent; do not reuse the same string for both unless you mean both behaviours.

Retention too short for jobId dedup. removeOnComplete: true removes the job immediately and ends its deduplication. Keep jobs as long as duplicates can arrive.

Flows. Deduplication applies per job; parents and children in flows have their own ids and rules.

Version differences. Deduplication modes arrived incrementally in 5.x; confirm the option names in your installed version's documentation.

FAQ

Is this the same as Sidekiq unique jobs? Similar in intent: both deduplicate at enqueue. See Sidekiq unique jobs and deduplication for the Ruby equivalent.

Can I deduplicate across queues? No — deduplication keys are per queue. For cross-queue uniqueness, use a Redis key or database constraint yourself, as in deduplicating jobs with Redis SET NX keys.

What happens to a debounced job that is already running? Once a job is active, a new request with the same deduplication id is no longer merged into it — depending on the mode and TTL, it is either ignored or creates the next pending job. For reindexing, that is what you want: the running job indexes the version it read, and the next one picks up later edits. Make the processor read current state at start, so a run that began just before the last edit is followed by one that includes it.

Can producers tell whether their add was deduplicated? Yes: the deduplicated event on QueueEvents reports it, and the returned job's id is the existing job's id rather than a new one. Log it at debug level; counting it as a metric is more useful than logging every occurrence.

Does deduplication cost much? One extra key per active deduplication id and a check in the add script — negligible next to the jobs it prevents.

Related