Consume and acknowledge messages¶
Take messages from a topic, process them, and acknowledge each one so that Narad never delivers it again.
Before you start: your user needs a consume grant that matches the topic. The same grant covers acks, extends and nacks.
Consume a message¶
curl -i -u "$AUTH" "$NARAD/v1/topics/orders/consume?wait=10s"
HTTP/1.1 200 OK
Content-Length: 180
Content-Type: application/json
Date: Mon, 28 Sep 2026 19:21:18 GMT
{
"topic": "orders",
"partition": 4,
"offset": 0,
"key": "customer-42",
"payload": {
"order_id": "ord_123",
"amount": 4999
},
"timestamp": 1790623277,
"receipt_handle": "4:0:3327358780392153864"
}
err := client.Consume(ctx, "orders", narad.HandlerFunc(
func(ctx context.Context, msg *narad.Message) error {
var order Order
if err := msg.Into(&order); err != nil {
return err
}
return process(ctx, order)
}))
Consume long-polls, renews the lease while the handler runs, acks when it returns nil and gives the message back when it returns an error. Go SDK has the details.
narad sub orders
[p4 @0] key=customer-42 03:06:49 {"order_id": "ord_123", "amount": 4999}
narad sub consumes and acks every message it prints, so it competes with your real consumers. Use narad sub orders --peek to watch without taking anything; see Replay messages.
What each part does:
wait=10sholds the request open for up to 10 seconds until a message arrives. If none does, the answer is204with no body. Withoutwait, the answer comes at once. Loop on it: an idle consumer costs one small request perwait.- The node caps
waitat the operator'shttp.max_consume_wait, 10 seconds by default. A longerwaitis not an error; the answer comes at the cap and carries anX-Narad-Wait-Clamped: 10sheader, so an early204has a visible reason. receipt_handleproves you hold the message. Treat it as opaque, send it back to ack, and never parse it.- From now on the message is hidden from other consumers for the topic's visibility timeout, 30 seconds by default. That is your lease: finish and ack within it, or extend it.
partition=Ntakes messages from that partition only. Every parameter and response field is in the HTTP API reference.
Run as many consumer processes as you like against the same topic. There are no consumer groups or partition assignments to configure: Narad hands each message to one consumer at a time.
No ordering guarantee
Narad does not guarantee delivery order. Messages with the same key usually arrive in the order they were produced, but redelivery, a broker restart and routing around an unreachable partition owner all reorder them. If you need a sequence, carry one in the payload.
The message lifecycle shows every state a message passes through, from produce to ack.
Delivery is at least once
Narad delivers a message until a consumer acks it, as long as the topic's retention still holds it, and it can deliver the same message more than once: after a lease expires, after a nack, or after a broker crash. Make every handler idempotent, for example by keying its work on an ID in the payload.
410 Gone.Acknowledge a message¶
curl -i -u "$AUTH" -X POST \
"$NARAD/v1/topics/orders/ack?receipt_handle=$HANDLE" \
-H "Content-Type: application/json"
HTTP/1.1 204 No Content
Date: Mon, 28 Sep 2026 19:21:27 GMT
err := msg.Ack(ctx)
Consume acks for you. Call Ack yourself only on a message you took with Receive.
$HANDLE is the receipt_handle from the consume response, here 4:0:3327358780392153864. A 204 means the message is settled and will not be delivered again. Acks may arrive in any order.
An ack that comes after the lease ran out, or for a message that is already acked, gets 410:
HTTP/1.1 410 Gone
Content-Length: 67
Content-Type: application/json
Date: Mon, 28 Sep 2026 19:21:27 GMT
{"error":"receipt handle no longer matches an active reservation"}
A 410 is not something to fix or retry. The message went back to the queue and may already be with another consumer, so your work on it may run twice; that is why handlers must be idempotent. An ack that gets 502 or 503 is different: it may not have reached the partition's owner, and you must retry it.
Extend a lease¶
A slow job can keep its message by renewing the lease instead of raising the topic's timeout for everyone:
curl -i -u "$AUTH" -X POST \
"$NARAD/v1/topics/orders/ack?receipt_handle=$HANDLE&extend=true" \
-H "Content-Type: application/json"
HTTP/1.1 204 No Content
Date: Mon, 28 Sep 2026 19:21:27 GMT
- A
204gives the message a full visibility timeout again, counted from now. Extend about every third of the timeout while you work, then ack with the same handle. - A
410means the lease already ran out. Stop working on the message; it belongs to someone else now. - The Go SDK's
Consumeextends for you while a handler runs.
Give a message back¶
When this consumer cannot handle a message right now, because it is shutting down or a dependency is down, hand it back at once instead of waiting out the lease:
curl -i -u "$AUTH" -X POST \
"$NARAD/v1/topics/orders/ack?receipt_handle=$HANDLE&extend=0" \
-H "Content-Type: application/json"
HTTP/1.1 204 No Content
Date: Mon, 28 Sep 2026 19:21:45 GMT
The message can be consumed again immediately, and a consumer waiting in a long poll gets it right away, with a new receipt handle. In the Go SDK, msg.Nack(ctx) does the same, and so does a Consume handler that returns an error.
A nack does not count anything. A message that always fails is nacked and delivered again forever; Handle retries and dead letters shows how to bound the attempts.
Consume a batch¶
New in v3.1.0.
curl -i -u "$AUTH" -H 'X-Narad-Client: curl' "$NARAD/v1/topics/orders/consume?wait=10s&max=10"
HTTP/1.1 200 OK
Content-Length: 373
Content-Type: application/json
Date: Mon, 28 Sep 2026 19:22:32 GMT
{
"messages": [
{
"topic": "orders",
"partition": 3,
"offset": 0,
"key": "customer-7",
"payload": {
"order_id": "ord_126",
"amount": 800
},
"timestamp": 1790623352,
"receipt_handle": "3:0:6999714737712221180"
},
{
"topic": "orders",
"partition": 4,
"offset": 4,
"key": "customer-42",
"payload": {
"order_id": "ord_125",
"amount": 1250
},
"timestamp": 1790623352,
"receipt_handle": "4:4:9069778768564137076"
}
]
}
- A batch consume must send an
X-Narad-Clientheader, with any value, or it gets400. The Go SDK and the CLI send it. A consume reserves messages, and a page on another site could send aGETwith an operator's cached Basic credentials; the header forces a CORS preflight, which Narad never approves. A single consume does not need it. maxis 1 to 100. Each message inmessagesis exactly what a single consume returns, with its own receipt handle and its own lease, and you settle each one on its own or in a batch. Anymax,max=1included, gets this shape; withoutmaxyou get one message, as before.- A batch is what is ready now. The request is never held to fill it, so expect fewer than
maxmessages. When nothing is ready it long-polls like a single consume and answers204ifwaitruns out. - A batch stops early when its messages add up to about 4 MiB, and a response carries at most 8 MiB. The first message always goes, whatever its size.
maxworks withpartition=Nbut not withoffset, which replays one message; that combination gets400.- A batch request counts as
maxrequests (or as the whole limit, whenmaxis larger) against your per-node limit on concurrent consumes, and each message counts against its partition'smax_in_flight_per_partition. - A node on v3.0.1 or earlier ignores
maxand answers one message in the single-message shape. Switch a client over once every node it can reach is upgraded.
Acknowledge a batch¶
New in v3.1.0.
Leave out the receipt_handle parameter and send the handles in a JSON body instead:
curl -i -u "$AUTH" -X POST "$NARAD/v1/topics/orders/ack" \
-H "Content-Type: application/json" \
-d '{"receipt_handles": ["3:0:6999714737712221180",
"4:4:9069778768564137076",
"4:0:3327358780392153864"]}'
HTTP/1.1 200 OK
Content-Length: 124
Content-Type: application/json
Date: Mon, 28 Sep 2026 19:22:37 GMT
{
"results": [
{
"status": 204
},
{
"status": 204
},
{
"status": 410,
"error": "receipt handle no longer matches an active reservation"
}
]
}
- Send 1 to 100 handles in a body of at most 64 KiB.
extend=trueorextend=0in the query applies to every handle. - The answer is
200with one result per handle, in request order. Eachstatusis what a single ack of that handle would have answered, anderrorexplains any status other than204. One stale handle never fails the others; here the third was already acked. - Retry only the handles whose result is
502or503. - A body that is malformed as a whole (invalid JSON, no handles, more than 100) gets
400, and a body over 64 KiB gets413. - A node on v3.0.1 or earlier answers
400withreceipt_handle required. As with batch consume, switch over once every node is upgraded.
Payload encoding¶
The payload field comes back in the form that carries the produced bytes exactly.
| You produced | You consume |
|---|---|
A JSON value, such as {"a":1}, [1,2], "hi" or 42 |
That JSON, byte for byte |
Text that is not JSON, such as hello world |
A JSON string: "payload":"hello world" |
| Binary data, such as an image or protobuf | A base64 string, flagged with "payload_encoding":"base64" |
Two consumes of the text and binary messages from Produce messages:
{
"topic": "orders",
"partition": 4,
"offset": 2,
"key": "customer-42",
"payload": "hello world",
"timestamp": 1790623346,
"receipt_handle": "4:2:3685415247247279671"
}
{
"topic": "orders",
"partition": 4,
"offset": 3,
"key": "customer-42",
"payload": "iVBORw0KGgo=",
"payload_encoding": "base64",
"timestamp": 1790623346,
"receipt_handle": "4:3:5733584771904418009"
}
The rule for consumers: if payload_encoding is "base64", decode the payload; otherwise use it as it is.
Keys follow the same idea. key is present only when the message has one, as a JSON string.
New in v3.1.0. A key that is not valid UTF-8 comes back base64-encoded with "key_encoding":"base64" beside it. During a rolling upgrade, a message whose partition owner still runs v3.0.1 comes back the old way, without key_encoding.
Flow control¶
Three limits keep one consumer from taking more than it can finish:
max_in_flight_per_partition(topic setting, default 1024): once that many messages from one partition are out and unacked, the partition hands out nothing more until acks arrive.max_acked_ahead_per_partition(topic setting, default 1024): the number of acks a partition holds for messages after the oldest unacked one. At the limit, the partition hands out only that oldest message until it is acked. Acks for messages you already hold are always accepted.- Concurrent consumes per identity, 1024 per node by default: past it, a consume gets
429. A long poll counts for as long as it waits.
The two topic settings can be changed at any time. Lag and in-flight counts per partition are in the metrics reference.
Next steps¶
- Handle retries and dead letters: bound retries, back off, and retry failed acks.
- Replay messages from an offset: read history without taking leases.
- Go SDK: produce and consume from Go: a consumer loop that manages leases for you.