A consumer group reads a queue at its own pace, with one cursor in each partition. Its workers take messages under a lease, and an ack moves the cursor past them. The lease is what keeps each entity’s events in order while many workers run at once, and the cursor is what lets a second group, or a replay, read the same messages without disturbing anyone.
await client
.queue('orders')
.group('billing')
.subscriptionMode('all')
.limit(1)
.each()
.consume(async (message) => {
console.log(message.data)
})consume long-polls for messages, calls your handler, acks the message when the handler returns
and nacks it (failed) when the handler throws, and keeps consuming either way. When the handler
acks inside a transaction, turn the automatic ack off with
.autoAck(false). A handler that throws is still nacked for whatever it had not acked, so its
partition comes back at once instead of waiting out the lease, and .onError() on consume hands
that failure to your own code instead (queen-mq 2.0; the 1.x SDK stopped the loop there). pop
hands you the leased messages and leaves the ack to you:
const messages = await client
.queue('orders')
.batch(10)
.wait(true)
.pop()curl -s 'http://localhost:6632/api/v1/pop/queue/orders?consumerGroup=billing&batch=10&wait=true&timeout=30000'A pop answers 200 with a leaseId and messages, each carrying its transactionId,
partitionId, offset, deliveryAttempt and data. A pop that finds nothing answers 204 with no
body.
Groups, cursors and leases
A group keeps one cursor per partition: the last offset it acked there. Two groups (billing and
analytics, say) read the same messages, each at its own pace, and neither can hold the other back.
Workers that share a group compete for its messages. A pop that names no group reads through the
queue’s shared group, __QUEUE_MODE__, which is the plain queue you would expect: every message
goes to one of the workers.
A pop leases a batch of one partition for one group, and no other worker of that group gets the partition until the lease ends. That single rule is why order holds: the next batch of a partition cannot be handed out while the previous one is still being worked on, and workers of the same group always work on different entities.
The lease lasts leaseSeconds from the pop when you send it, else the queue’s leaseTime (60 s
for a queue a push created). A worker that needs longer renews it with
POST /api/v1/lease/:leaseId/extend and {"seconds": 60} (in JS, queen.renew(message), or
.renewLease(true, 10000) on consume); a lease that has already expired cannot be renewed. When
a lease expires, the partition is free again and its batch is delivered to the next pop, from the
cursor, with deliveryAttempt one higher.
Ack, nack and the retry budget
curl -s -X POST http://localhost:6632/api/v1/ack -H 'content-type: application/json' -d '{
"transactionId": "order-9137-created", "partitionId": "4", "consumerGroup": "billing",
"leaseId": "01a0fcc5-7656-7000-8144-b8981abf548a", "status": "completed" }'| Status | Effect |
|---|---|
completed |
Moves the cursor to this message. The earlier messages of the batch are done with it. |
failed |
Releases the lease and spends one retry: the message comes back from this point, with deliveryAttempt one higher. With retryLimit spent, files it in the dead-letter queue, or drops it when the DLQ is off. |
retry |
Releases the lease and redelivers from this message, without spending a retry. |
dlq |
Files the message in the dead-letter queue now and moves on. |
In an ack of several messages (POST /api/v1/ack/batch with
{"acknowledgments": [...], "consumerGroup": "billing"}), the cursor stops at the lowest message
that will come back, a failed with retries left or a retry, and a completed above it is
delivered again too. The retry budget is one counter per partition and group, and only an
explicit failed spends it. Every ack answers 200, even when an item is refused, so read each
item’s success.
A dead letter keeps its payload, the error the worker sent, the retry count and the group that gave up on it, and it stays until you replay or delete it:
curl -s 'http://localhost:6632/api/v1/dlq?queue=orders'
curl -s -X POST http://localhost:6632/api/v1/dlq/<dead-letter-id>/replayA replay pushes a copy to the end of the partition, under the transactionId
dlq:<dead-letter-id>, so every group that reads that partition sees it, not only the group that
gave up.
Where a new group starts
A group’s starting point is decided on its first pop and stored, so sending a different mode later changes nothing.
subscriptionMode=new, the default (DEFAULT_SUBSCRIPTION_MODE): only messages pushed after the group’s first pop.subscriptionMode=all: everything retained, from the oldest message.subscriptionFrom=<ISO time>: messages pushed from that time on.
A pop with no group ignores these and starts from the oldest retained message. A group that has
already read moves with a seek (POST /api/v1/consumer-groups/:group/queues/:queue/seek, to a
timestamp or to the end), which also releases any live lease:
// Move the audit group's cursor back one hour. The seek also releases
// any live lease, so an in-flight batch is abandoned, not acked.
await client.admin.seekConsumerGroup('audit', 'orders', {
timestamp: new Date(Date.now() - 3600 * 1000).toISOString(),
})Pop options
A pop takes its options in the query string. The SDKs expose the same ones as builder methods.
wait=true&timeout=30000long-polls: a queue or partition pop waits until a message arrives or the timeout passes (30 s by default, fromDEFAULT_TIMEOUT, and at most 60 s per request).partitions=N(1 to 64) claims up to N partitions in one pop.batchis then the budget of the whole pop, and every claimed partition shares oneleaseId.autopilot=trueon a queue pop lets the broker choosebatch(100 to 1,000, from how fast the group acks) andpartitions(the ready partitions, shared among the waiting pops). A value you send is never changed. The JS SDK turns autopilot on by default;.autopilot(false)turns it off.GET /api/v1/pop?namespace=billing&task=invoicesis a discovery pop, across every queue with those labels (see queue options). It does not long-poll, and its lease is 60 s unless the pop sendsleaseSeconds.conflation=trueon a group delivers only the newest message of a partition and moves the cursor past the older ones, which suits command-style queues where only the latest input matters. The group’s first pop decides it; it needs a group and refusesautoAck, and an empty pop of a conflating group answers 200 with no messages.autoAck=truecommits the cursor when the messages are handed out. That is at-most-once delivery: a crash after the pop loses the batch.
Limits
- Delivery is at-least-once: an expired lease, a nack or a crash delivers a message again. Make the step safe to repeat with a transaction and dedup.
- Parallelism inside a group is bounded by the partitions that have work. Workers beyond that number wait.
- An expired lease never spends the retry budget. A handler that throws is nacked, which spends it, but a message that kills its worker’s process every time comes back forever and never reaches the DLQ.
- An ack without a
leaseIdis not fenced. The broker accepts it from a worker whose lease expired, and it can complete, and release, the batch another worker holds now. Send theleaseIdthe pop returned (the SDKs do). - Lease exclusivity across a leader change needs the nodes’ clocks within
QUEEN_RAFT_MAX_CLOCK_SKEW_MS(500 ms) of each other. - A new group skips the backlog by default. Use
subscriptionMode('all')to read what was pushed before it existed.