Skip to content

Dedup and once

A transactionId makes a push idempotent inside its partition's dedup window; once makes a whole step happen at most once.

Updated View as Markdown

Every message carries a transactionId. Push the same id to the same partition again, inside the queue’s dedup window, and the broker writes nothing and answers with the original message. That makes a retried push harmless, and you will retry: a push that timed out may or may not have been stored, and sending it again is the only way to find out. When a whole step must not run twice, once does the same for the step.

clients/client-js/test-v2/docs.jsjs
const first = await client
  .queue('payments')
  .partition('customer-42')
  .push([{ transactionId: 'order-9137-paid', data: { orderId: 9137, amount: 99.5 } }])

const retry = await client
  .queue('payments')
  .partition('customer-42')
  .push([{ transactionId: 'order-9137-paid', data: { orderId: 9137, amount: 99.5 } }])
// retry[0].status is 'duplicate': the second push wrote nothing
// and answers with the first message's id.

Over HTTP the retry answers HTTP 201 with "status": "duplicate", the original message_id and the original offset, and nothing is appended.

How the check works

The leader checks each pushed id against the partition’s index of recent ids when it plans the push, in the same step that assigns offsets, so the answer is exact. The index is part of the replicated state: every node records the ids as it applies an entry, so after a failover the new leader answers a retry the same way, with the original offset. That case was Jepsen-tested.

A producer pushes a message with transactionId order-9137-v3. The leader checks the partition's index of recent ids, stores the message at offset 41 and answers 201 queued, but the answer is lost and the producer sees a timeout. The leader then fails over. The producer pushes the same id again, the new leader finds it in the replicated index and answers 201 duplicate with the original offset 41, appending nothing.producerleadernode 1new leadernode 2push order-9137-v3new id: offset 41201 queued, losttimeout: stored or not?node 1 fails, node 2 leadspush order-9137-v3 againid seen: nothing appended201 duplicate, offset 41
The index of recent ids is replicated state, so the answer to a retry does not depend on which node leads. The window is the queue's dedupWindowSeconds, an hour by default.

The scope is one partition. The same id in another partition, or another queue, is a different message, which is why every retry of a message must go to the same partition. The window is the queue’s dedupWindowSeconds, 3,600 seconds by default; 0 turns the check off. Inside one request, an id repeated for the same queue and partition is stored once, whatever the window, and the repeats answer duplicate. A push without an id gets the message id as its transactionId, which is new every time, so nothing is deduplicated (the JS SDK mints a UUIDv7 when you pass none, to the same effect).

Queues created for a topic of the Kafka facade are the exception: since 2.0.0-beta.5 they get a window of 0. A Kafka record carries no transactionId, so the broker mints a fresh one for each record and the check could only ever answer no, and Kafka’s idempotent producer suppresses its own retries with its sequence numbers. A queue that already existed keeps its own setting.

Choose deterministic ids

Derive the id from the work, never from the attempt: an upstream event id, a primary key with a version (order-9137-v3), or the input’s id in a pipeline (invoice-${message.transactionId}). A fresh UUID per attempt deduplicates nothing.

A duplicate inside a transaction rolls back everything

When a push inside a transaction carries an id its partition already stores, the whole transaction rolls back with reason: "duplicate": the ack, the KV writes and the timers are not written either. That is what makes a retried transaction safe. If the first attempt committed, the retry finds its push ids and writes nothing; if it didn’t, the retry commits. (An id repeated within the same transaction is different: it is stored once, answered with "duplicate": true, and the transaction commits.)

once: a step at most once

A redelivered event carries the same transactionId, but your step would run again. once writes a marker in the step’s transaction and refuses the whole transaction when the marker already exists:

const res = await queen.transaction()
  .ack(message, 'completed', { consumerGroup: 'mailer' })
  .queue('emails').push([{ data: mail }])
  .once('welcome', userId, { ttl: '7d' })
  .commit()
if (res.success === false) await queen.ack(message, 'completed', { group: 'mailer' })  // already done

once(ns, key, opts) is a KV putIfAbsent with required: true. A lost gate answers HTTP 200 with reason: "kv_precondition", and the JS SDK returns that answer from commit() without throwing. The ack rolled back too, which is why the example acks on its own. Outside a transaction, queen.kv.once(ns, key, { ttl }) returns { won }.

A transaction checks its pushes before its KV operations. If it also pushes with a deterministic id, a repeat rolls back with duplicate first, and the JS SDK throws that one; handle both, or let once be the only gate, as above.

You want Use
a retried push not to store twice a deterministic transactionId
a redelivered event not to run its step twice once in the step’s transaction
an external call not to run twice an idempotency key the other system honours

Limits

  • The window is bounded: outside dedupWindowSeconds, the same id is a new message.
  • Dedup is per partition, so every retry of a message has to go to the same partition.
  • A once marker lives as long as you say (ttl, or forever: true). After it expires, a redelivery runs the step again, so make the TTL longer than any redelivery can arrive.
  • Dedup does not cover a side effect outside Queen. See exactly-once.
  • A transactionId is at most 65,535 bytes.

Next

Navigation

Type to search…

↑↓ navigate↵ selectEsc close