Skip to content

Transactions

One POST /api/v1/transaction writes a whole step (acks, pushes, KV writes, timers) as one log entry: all of it applies or none of it does.

Updated View as Markdown

A transaction is one call that writes one step. It acks what the worker took, pushes the next events, writes KV state and schedules or cancels timers, as a single entry of the broker’s log, and either every part applies or none does. Use it whenever a worker’s output must not come apart from its input, which in practice is almost always.

await queen.queue('payments').group('billing').autoAck(false).each()
  .consume(async (message) => {
    const { orderId, customerId, paymentId } = message.data
    try {
      const res = await queen.transaction()
        .ack(message, 'completed', { consumerGroup: 'billing' })       // progress
        .kv.put('orders', `order-${orderId}`, { status: 'paid' }, { ttl: '30d' }) // state
        .queue('receipts').partition(customerId)
          .push([{ data: { orderId, paymentId } }])                    // events
        .timer('reminders').key(`order-${orderId}`).delay('24h')
          .payload({ orderId }).schedule()                             // time
        .once('payments', paymentId, { ttl: '30d' })                   // identity
        .commit()
      if (res.success === false) {
        // `once` found its marker: this payment was handled before.
        await queen.ack(message, 'completed', { group: 'billing' })
      }
    } catch (err) {
      // rejected_ack: the lease ran out, and the message goes to the next worker as it is.
      if (err.reason !== 'rejected_ack') {
        await queen.ack(message, 'failed', { group: 'billing', error: err.message })
      }
    }
  })

The second delivery of the same payment (a provider retrying its webhook, a lease that expired after the commit) finds the once marker, and the whole transaction rolls back, receipt and timer included. The handler then acks the message on its own, because the ack rolled back with the rest.

Over HTTP the same call is POST /api/v1/transaction, with the acks and pushes in operations and the KV and timer operations in the top-level arrays kv and timers; the model shows a whole body, and the reference every field. A commit answers HTTP 200 with success: true and one result per ack, pushed item, KV operation and timer operation. A rollback answers HTTP 200 too, with success: false, a reason, an error and an empty results, and it has written nothing.

What can go in

Part Where What it can hold
Acks operations, type: "ack" Messages leased from any partitions, queues and consumer groups of the tenant. Each names transactionId, partitionId, consumerGroup (default __QUEUE_MODE__) and its leaseId.
Pushes operations, type: "push" Messages to any queues and partitions. A queue or partition that does not exist is created, as by a push.
KV kv put, putIfAbsent, delete, incr, get, getMany. No getPrefix. See KV.
Timers timers schedule (an upsert) and cancel. See timers.
One call to POST /api/v1/transaction carries four parts: acks of leased messages, pushes of the next events, KV operations and timer operations. The broker checks them in that order, acks first, then the pushes and their dedup, then KV, then timers. If every check passes, all four parts are written as one log entry and apply together. The first refusal decides the reason (rejected_ack, duplicate, kv_precondition and others), and the call rolls back with nothing written.acksleased messagespushesnext eventsKVstatetimersschedule, cancelone transactionchecked: acks, pushes, KV, timersone log entryevery part appliesrolled backnothing writtenall passfirst refusal
Four parts, one entry, one verdict. The checks run in a fixed order, so the first refusal names the reason. Source: server/src/rsm/planner/txn.rs, server/src/rsm/consume/txn.rs

Every ack belongs to a consumer group. The JS SDK takes it from the message, which every pop answers with (message.consumerGroup). Over HTTP, name it on every ack operation: an ack without one is judged in queue mode (__QUEUE_MODE__), finds no lease there, and the whole call rolls back. The entry is atomic across a crash, too: it is written into the log of every queue it touches, with the number of copies, and recovery replays it only when every copy is present, so a transaction is never partly visible.

The lease is the fence

When an ack carries the leaseId of the pop that delivered its message, that lease fences the whole call. If the lease has expired, or another worker holds the partition now, the ack is refused and the transaction rolls back with reason: "rejected_ack", KV writes and timers included. A compare-and-swap cannot give you this: an expect on a version that still matches succeeds even for a worker that lost its message minutes ago.

Worker A pops a batch and gets lease L1, then takes longer than the lease. The lease expires, and worker B's pop gets the same batch under lease L2. When worker A commits its transaction (the ack under L1, a KV write and a push), the broker refuses it with rejected_ack and writes none of it. Worker B's transaction, under L2, commits.worker AQueenworker Bpopbatch, lease L1slow: L1 expirespopsame batch, lease L2ack (L1) + KV + pushrejected_ack, nothing writtenack (L2) + KV + pushsuccess: one entry
The ack carries the lease, and the lease decides for the whole call. A worker that lost its message cannot write its state late, which a compare-and-swap on a version alone would allow.

The lease has to travel with the ack. Over HTTP, put leaseId on each ack operation, or name exactly one lease in requiredLeases and the acks without their own leaseId take it; an ack that ends up with no lease is not checked and still moves the cursor. The SDKs put each message’s lease on its ack.

once, the gate

.once(ns, key, { ttl }) writes a KV marker with putIfAbsent and required: true. When the marker already exists, the transaction rolls back and answers success: false, reason: "kv_precondition", kvReason: "exists", and the failedIndex of the gate. The JS SDK returns that verdict from commit() without throwing, because a second event for the same business fact (a provider retrying its webhook, a redelivery after a crash) is something your handler expects and answers with an ack. Key the gate on that fact (paymentId), and give the marker a TTL longer than any redelivery can arrive.

Why a transaction rolls back

reason Cause JS commit()
kv_precondition A KV operation with required: true lost its precondition (once found its marker, an expect missed). returns the verdict
duplicate A pushed transactionId is already in its partition’s dedup window. throws, err.reason
rejected_ack An ack’s lease has expired or moved, or the message it names is no longer leased to this worker. throws
reserved Another transaction or a checkpoint holds that cursor at this moment. Retry. throws
bad_request A malformed operation. Some shapes answer HTTP 400 instead. throws
too_large One push group, KV call or timers part is over the entry limit. throws

The broker checks the acks first, then the pushes and their dedup, then KV, then timers, and the first refusal decides the reason. Errors lists every code. A 503 is not a rollback: retry, timeout and no_leader mean the outcome is not known (the call may still commit), and they carry Retry-After.

One tenant, one raft group

A transaction carries the tenant of its request, and each tenant lives in exactly one raft group (QUEEN_RAFT_GROUPS, placed by a hash of the tenant name or by QUEEN_TENANT_GROUPS). So every transaction is one entry in one log, ordered by one leader, with no coordinator and no two-phase commit. That is the design decision the rest of this page leans on, and its cost is that a transaction cannot span tenants.

Where exactly-once stops

Answers get lost. The SDKs send a call again on their own when it fails with a 5xx or times out, and you may too, so a step can reach the broker twice. A retry of a step that already committed is refused whole when the step carries a gate: a once marker, or a push whose deterministic transactionId the first attempt stored (see dedup). The ack alone is not always a gate. If the first attempt released the lease, the repeated ack is refused, but while the same lease still covers other messages of the batch, re-acking an acked message is a no-op, and a retry without a gate writes its KV operations and timers again (an incr counts twice). So give every step a gate.

In JavaScript, a handler under .autoAck(false) that throws is nacked for whatever it had not acked, and the loop goes on, so an error before the commit spends one retry and the event comes back. Catch the errors you can answer better yourself, as the example above does with a lost lease.

A call your worker makes to another system, such as a payment provider, is outside the commit. Give that system an idempotency key derived from the message. See exactly-once.

Limits

  • A transaction is a single call, not an interactive session: there is no BEGIN, read, decide, COMMIT. A get inside a transaction is answered after the entry has applied, so it cannot be a condition; use expect with required: true.
  • One transaction belongs to one tenant.
  • KV inside a transaction takes at most 64 operations over 256 keys, one write per key, and no getPrefix.
  • One call holds one timer operation per queue and timerKey.
  • One body is at most 64 MiB (QUEEN_MAX_BODY_BYTES), one planned entry at most 96 MiB (QUEEN_RAFT_ENTRY_MAX_BYTES).
  • Throughput has a ceiling per raft group. In our 2026-10-01 runs on three 16-vCPU nodes, 2.0.0-beta.2 committed up to about 7,200 transactions a second of 10 messages each on 200 partitions (8,900 in a tuned run), with a commit p99 of 4 ms at 900 a second (benchmark-queen/2026-09-30-kafka-pulsar/runs/txn/; see transactions). Each raft group is a whole state machine with its own leader, so tenants in other groups have a ceiling of their own.

Next

Navigation

Type to search…

↑↓ navigate↵ selectEsc close