Skip to content

Consuming

Consumer groups read partitions through leases: one cursor and one leased batch per partition and group, acks that move the cursor, a retry budget and a dead-letter queue.

Updated View as Markdown

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.

clients/client-js/test-v2/docs.jsjs
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:

clients/client-js/test-v2/docs.jsjs
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.

Partition 4 of the queue orders holds offsets 0 to 7, stored once. Consumer group billing has acked up to offset 5 and holds a lease on offsets 6 and 7. Group analytics has acked up to offset 2. Each group reads the same messages at its own pace.billing's lease01234567billing acked up to 5analytics its own pace
One copy of the messages and a cursor per group. The lease billing holds on 6 and 7 holds back nobody in analytics.

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.
What happens to a delivered message. Acked completed, it moves the group's cursor past it. Acked failed with retries left, it releases the lease, spends one retry and comes back from that point with deliveryAttempt one higher. Acked failed with the retry budget spent, or acked dlq, it is filed in the dead-letter queue with its payload, error, retry count and group. A replay pushes a copy to the end of the partition under the transactionId dlq:<id>.partitionat the cursorworkergroup billingcursor moveson to the nextdead-letter queuepayload, error, retriesdeliveredfailed: retrycompletedfailed, budget spent, or dlqreplay: a copy at the end
Only an explicit failed spends the budget; a lease that expires delivers again without spending it. A replay is a new message, so every group reading the partition sees it.

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>/replay

A 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 partition with retained offsets 0 to 7 and where a new group starts under each mode. With subscriptionMode=all it starts at the oldest retained message, before offset 0. With subscriptionFrom set to a time it starts at the first message pushed from then on, offset 4 here. With subscriptionMode=new, the default, it starts after offset 7 and reads only what is pushed after its first pop.01234567all the oldest retainedsubscriptionFrom from a timenew the default
The starting point is decided on the group's first pop and stored. After that, only a seek moves it.

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:

clients/client-js/test-v2/docs.jsjs
// 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=30000 long-polls: a queue or partition pop waits until a message arrives or the timeout passes (30 s by default, from DEFAULT_TIMEOUT, and at most 60 s per request).
  • partitions=N (1 to 64) claims up to N partitions in one pop. batch is then the budget of the whole pop, and every claimed partition shares one leaseId.
  • autopilot=true on a queue pop lets the broker choose batch (100 to 1,000, from how fast the group acks) and partitions (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=invoices is 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 sends leaseSeconds.
  • conflation=true on 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 refuses autoAck, and an empty pop of a conflating group answers 200 with no messages.
  • autoAck=true commits 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 leaseId is 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 the leaseId the 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.

Next

Navigation

Type to search…

↑↓ navigate↵ selectEsc close