Skip to content

Delivery contract

Learn what Narad promises about delivery, durability, ordering and availability, and what each kind of failure does to your messages.

In short

  • Delivery is at least once. A message is delivered until a consumer acks it, sometimes more than once, so every handler must be idempotent.
  • A 202 means the node that answered fsynced the message. It survives any crash, restart or power loss from then on.
  • Narad stores each partition once. Losing a node's volume loses the partitions on it, unless you keep a replica child or volume snapshots.
  • There is no ordering guarantee. Carry a sequence in the payload if you need one.
  • Produce and consume keep working while nodes fail. Topic and user changes need a Raft quorum.

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.

A consumer takes a lease on each message it receives. The lease lasts the topic's visibility timeout (visibility_timeout_ms, 30 seconds by default), and the response carries a receipt handle that names the lease. An ack with that handle settles the message for good. Anything that ends a lease without an ack puts the message back in the queue.

flowchart LR
    accTitle: The life of one message under at-least-once delivery
    accDescr: A message accepted with 202 is delivered. If it is acked it is settled and never delivered again. If its lease ends without an ack it is delivered again, possibly as a duplicate. If it is never acked before its retention runs out, retention deletes it.
    A["202 Accepted"] --> B{delivered and acked?}
    B -->|yes| C[settled for good]
    B -->|"lease ended, no ack"| D["delivered again<br/>(maybe a duplicate)"]
    D --> B
    B -->|"never acked within retention"| E[deleted by retention]

These events deliver a message again:

  • The lease lapses. The consumer crashed, hung or worked past the visibility timeout. The message goes to the next consumer, and the late ack is answered 410.
  • A consumer nacks it. A nack ends the lease at once.
  • The broker process crashes. Leases live only in the memory of the partition's owner. When the node comes back, every message that was leased on its partitions is delivered again straight away, even if the first consumer is still working on it. Acks reach the disk in batches, so the messages acked in about the last 100 ms come back too.
  • A node loses power, its kernel crashes, or its machine is lost with the volume intact. As for a crash, but acks of about the last 1.1 s come back by default: the durability interval storage.consumer_offset_commit_interval_ms (1 s, from v3.1.0), plus about 100 ms and the time a sync takes. v3.0.1 brings back about 0.1 s plus its sync time.
  • A produce is retried after a timeout. The first attempt may have been accepted, so the retry can store a second copy.
  • A node commits accepted messages again after a crash. Its dispatcher keeps in memory which messages above its checkpoint are already committed, so a crash commits those again at new offsets. After a power loss the checkpoint can also be up to 250 ms old (from v3.1.0; v3.0.1 syncs it on every store).
  • A partition's owner crashes after storing a commit but before answering it. The node that accepted the messages cannot tell that the commit landed, so it sends them again: to the same partition once the owner is back, or, after 3 s of failures, to a live sibling partition. Each message is then stored twice, and each copy is delivered and leased on its own, so two consumers can hold and ack the two copies at the same time.
  • A partition moves. Leases still out when a rebalance or decommission hands the partition to a new owner are delivered again by that owner.
  • A remote child resends a request (from v3.2.0). When a request to the remote times out or fails without an answer, the cursor sends its records again, so the copy on the remote can hold a record twice, each delivered on its own there (What a remote child promises).

A graceful stop writes every ack to disk first, so a rolling restart delivers no acked message again. It does forget the leases of the stopping node, as a crash does.

None of these loses a message. The one way an unacked message leaves without being delivered is retention. The patterns for idempotent handlers are in Handle retries and dead letters.

What a 202 means

A 202 Accepted means the node that answered has written your message to its write-ahead log and fsynced it. From then on a crash, restart or power loss of any node does not lose it. Consumers do not see it yet: the node hands it to the partition's owner in the background, usually within milliseconds. A produce that times out may or may not have been accepted, so retrying it can create a duplicate.

A produce is written to the ingress WAL of the node that received it, and the 202 goes out once the group-commit fsync that covers it returns. That node's dispatcher then commits the message to the owner of its partition. The owner fsyncs it and reads it back with its checksum verified, and only then moves the high watermark over it so consumers can see it. The ingress WAL keeps its copy until the owner confirms that commit, so a message that got a 202 is always in at least one verified place.

A produce that gets no answer is ambiguous, and so is one answered 500. When the ingress WAL's disk fails in the middle of a write, the records written before the failure survive the restart and are delivered (when the WAL's disk fails). A batch produce (v3.1.0) gets one 202 for all of its messages, and a batch that fails or times out may have been accepted in part.

The whole path, stage by stage, is in Produce path.

One copy per partition

Narad stores each partition once, on the volume of the node that owns it. A crash, restart or power loss loses nothing that got a 202; losing a volume loses the partitions on it. For a second copy, add a replica child or take volume snapshots, as Back up and replicate topics shows.

Narad has no replication subsystem for message data. Each partition is a directory on its owner's volume, and every commit is fsynced and verified there before it becomes visible. So process crashes, restarts and power loss lose nothing that got a 202. A destroyed volume loses the messages, consumer positions and fan-out cursors stored on it.

Cluster metadata is different. Topics, users, grants, schemas and partition assignments live in the Raft-replicated metastore, which every node holds in full. It survives the loss of any minority of volumes.

Two tools protect message data against a lost volume, and Back up and replicate topics covers both:

  • A replica child is an asynchronous full copy of a topic, placed on nodes other than the parent's partitions when it is created. It trails the parent by the fan-out lag.
  • Volume snapshots give each node a restore point.

New in v3.2.0: a remote child also keeps an asynchronous copy of a topic on another Narad cluster, which survives the loss of this whole cluster or its region. It trails the parent by the link's lag (Set up disaster recovery).

Run Narad on storage you trust, such as cloud persistent volumes or RAID.

Ordering

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.

Messages with the same key go to one partition by the hash of the key, and in steady state they tend to arrive in the order they were produced. Five mechanisms reorder them on purpose, and a design must assume all five:

  1. Redelivery. A message whose lease lapsed, or that was nacked, comes back after newer messages were consumed.
  2. A stranded lease holds its partition. A message leased by a consumer that never comes back is delivered again only when its visibility timeout ends. Until then its partition may have nothing else to serve, because everything above it is already acked. So after an outage a partition can go quiet for up to one visibility timeout and then deliver the rest (the mechanism).
  3. A broker restart. Acks reach the disk in batches, so after a crash the messages acked in the last moments are delivered again, after newer ones.
  4. Each node dispatches on its own. A message is committed to its partition by the dispatcher of the node that accepted it. Two messages with the same key that reach two different nodes can land on their partition in either order.
  5. Reroute around an unavailable owner. When a partition's owner is marked dead, or its commits keep failing for 3 seconds, the accepting node commits those messages to a live sibling partition instead (the reroute). The key-to-partition mapping moves for them.
Same key, two nodes: the order is not kept Your service sends m1 with key customer-42 to narad-0, which answers 202. After that 202 it sends m2 with the same key to narad-1, which also answers 202. The key maps both to partition orders/1, owned by narad-2. Each node commits what it accepted with its own dispatcher, on its own clock. 3: narad-1's dispatcher commits m2 first, at offset 7. 4: narad-0's dispatcher commits m1 after it, at offset 8. So m1, sent first, lands second. Your service narad-0 dispatcher narad-1 dispatcher narad-2 orders/1 6 m2 7 m1 8 1 m1, key=customer-42 202 2 then m2, same key 202 m1 4 second: offset 8 3 first: offset 7 each node lingers and batches on its own clock, up to 50 ms Same key, two nodes: the order is not kept Your service sends m1 with key customer-42 to narad-0, which answers 202. After that 202 it sends m2 with the same key to narad-1, which also answers 202. The key maps both to partition orders/1, owned by narad-2. Each node commits what it accepted with its own dispatcher, on its own clock. 3: narad-1's dispatcher commits m2 first, at offset 7. 4: narad-0's dispatcher commits m1 after it, at offset 8. So m1, sent first, lands second. Your service narad-0 dispatcher narad-1 dispatcher 1 2 m1 3 4 narad-2 orders/1 6 m2 7 m1 8 1 m1, key=customer-42 to narad-0: 202 2 then m2, same key, to narad-1: 202 3 narad-1's dispatcher commits m2 first, at offset 7 4 narad-0's dispatcher commits m1 second, at offset 8 each node lingers and batches on its own clock, up to 50 ms
Two nodes, two dispatchers, two clocks: m1 was sent first and lands second, at offset 8.

If you need a sequence, carry one in the payload and order on your side. Handlers that are idempotent on an ID in the payload absorb duplicates and reordering together.

Availability

Ordering was traded for availability. In CAP terms, Narad's data plane is AP and its control plane is CP.

  • Produce works while any node lives. Any live node accepts a produce with a local fsync: no leader election, no quorum, no coordination on the request path. Delivery to the owner happens afterwards and routes around dead nodes. When a majority of nodes is down, a survivor still accepts a produce sent to it directly. It reports not ready while it has no Raft leader, though, so a load balancer stops sending it traffic until quorum returns. Losing a minority of nodes never stops produces through the load balancer.
  • Consume works for every partition whose owner is alive. New messages reroute to live owners, so fresh messages stay consumable during an outage. Messages already stored on a dead node wait for it to return. Meanwhile a consume pinned to one of its partitions is answered 503, or 502 until the node is marked dead.
  • Topic, user and grant changes go through Raft and need a quorum of nodes. Without one they are answered 503. Data flows through one node; administration waits for a majority.
What still works with one node down, and with two Two views of a three-node cluster behind a load balancer. With one of three nodes down, narad-2, the load balancer sends traffic to narad-0 and narad-1. A produce, ord_123, is answered 202 by any live node. Consumes work, except that messages stored on narad-2 wait for it to return, and a consume or an ack pinned to one of its partitions is answered 502 while the node does not answer, then 503 once the controller marks it dead, about 30 seconds after its last heartbeat. Topic and user changes work, because Raft still has a quorum, 2 of 3. With two of three down, narad-1 and narad-2, narad-0 has no Raft leader and reports not ready, so the load balancer routes nothing to it. A produce sent to narad-0 directly is still answered 202, and a consume sent to it directly is served from what it stores. Topic and user changes are answered 503: there is no quorum. One of three nodes down load balancer narad-0 narad-1 narad-2 ord_123 narad-0, narad-1: ready produce: 202 on any live node consume: yes, except messages stored on narad-2: they wait; pinned to them: 502, then 503 once it is marked dead topics and users: yes, Raft quorum 2 of 3 Two of three nodes down load balancer narad-0 narad-1 narad-2 routes nothing narad-0: not ready, no Raft leader produce sent to narad-0 directly: still 202 consume sent to narad-0 directly: what it stores topics and users: 503 no Raft quorum What still works with one node down, and with two Two views of a three-node cluster behind a load balancer. With one of three nodes down, narad-2, the load balancer sends traffic to narad-0 and narad-1. A produce, ord_123, is answered 202 by any live node. Consumes work, except that messages stored on narad-2 wait for it to return, and a consume or an ack pinned to one of its partitions is answered 502 while the node does not answer, then 503 once the controller marks it dead, about 30 seconds after its last heartbeat. Topic and user changes work, because Raft still has a quorum, 2 of 3. With two of three down, narad-1 and narad-2, narad-0 has no Raft leader and reports not ready, so the load balancer routes nothing to it. A produce sent to narad-0 directly is still answered 202, and a consume sent to it directly is served from what it stores. Topic and user changes are answered 503: there is no quorum. One of three nodes down load balancer narad-0 narad-1 narad-2 ord_123 narad-0, narad-1: ready produce: 202 on any live node consume: yes, except messages stored on narad-2: they wait; pinned to them: 502, then 503 once it is marked dead topics and users: yes, Raft quorum 2 of 3 Two of three nodes down load balancer narad-0 narad-1 narad-2 routes nothing narad-0: not ready, no Raft leader produce sent to narad-0 directly: still 202 consume sent to narad-0 directly: what it stores topics and users: 503 no Raft quorum
Losing a minority stops nothing that goes through the load balancer. Losing a majority stops the control plane and the load balancer, but a produce sent straight to a survivor still gets 202.

Retention

Retention is the one way an unacked message leaves without being delivered. A topic's retention_ms (at least 1 hour, or 0 to keep messages forever; when unset, the operator's default: 7 days for the binary, 12 hours for a cluster installed with the Helm chart) is a floor: a message lives at least that long after it was written. It usually lives somewhat longer, because deletion works on whole segments of up to 64 MiB, and it is gone within about twice the retention age of its write.

When a consumer falls so far behind that its next message has been deleted, the partition's committed frontier jumps to the oldest message still retained, and the owner logs consumer frontier fell behind retention; skipped to oldest retained offset. A fan-out child that falls behind its parent's retention skips the lost records too, and counts them in narad_fanout_child_dropped_messages. The parent of a delay child must keep messages for at least the delay plus one hour, which keeps a healthy delay child clear of this.

How segments age out is in Storage engine.

Timing

  • Produce to consumable: typically single-digit milliseconds. Under load, tens of milliseconds while the dispatcher gathers messages into larger commits.
  • Redelivery after a lapsed lease: one visibility timeout after the lease was taken or last extended. A sweep releases expired leases every second, and every consume of the partition releases them too.
  • Delay children: never early on the clock of the node that owns the parent partition, and usually within a second after the delay has passed. Failures can make them later.
  • Ack persistence: about 100 ms to the operating system and about 1 s to disk by default, as listed under At least once.

Failure matrix

Each row is one event: what clients see while it lasts, what it can lose, and where the procedure is.

Event Producers see Consumers see What can be lost What to do
A consumer crashes holding a lease Nothing The message again after its visibility timeout Nothing Make handlers idempotent
A handler outlives the visibility timeout Nothing Another consumer gets the message; the late ack gets 410 Nothing, but the work may run twice Extend the lease while working
An ack is answered 503 Nothing It did not land Nothing Retry the ack with backoff; a 410 on the retry means the lease is gone
An ack is lost, or answered 502 Nothing Unknown whether it landed; a retry of one that landed gets 410 Nothing Retry the ack with backoff
A produce times out Unknown whether it was accepted A retried produce may arrive twice Nothing, once retried Retry, per the retry rules
The broker process crashes Requests in flight to that node fail; other nodes accept Its partitions wait. On its return, its leased messages and about 100 ms of acks come back at once, and messages committed to it just before the crash can arrive twice Nothing that got a 202 Nothing, with idempotent handlers
A node loses power, volume intact As for a crash As for a crash, plus about 1.1 s of acks, and about 250 ms of dispatched messages twice Nothing that got a 202 To shorten the ack window, lower the durability interval
A node is down 202 as usual; its partitions' messages go to other partitions after about 3 s Its stored messages wait; pinned consumes fail with 502 or 503. After it returns, a partition can stay quiet for one visibility timeout Nothing Bring it back, then see quiet after an outage
A node's volume is destroyed The node rejoins empty; produce continues Its messages never arrive Every message, consumer position and fan-out cursor on it Restore a snapshot, or move consumers to a replica child
The ingress WAL's disk fails 500 from that node until it restarts, even after space is freed No change Nothing that got a 202; a produce answered 500 may still arrive Alert on narad_ingress_wal_failed (v3.1.0), then fix and restart
Raft quorum is lost Survivors report not ready, so the load balancer stops routing to them; topic and user changes get 503 The load balancer stops routing to the survivors Nothing Bring nodes back; do not restart the survivor
A partition moves (rebalance or decommission) 202 as usual; commits to it pause for the freeze, usually milliseconds No new messages from it during the freeze; leases unacked at the handover come back from the new owner Nothing Nothing; follow it with narad cluster moves
A move's source dies mid-move As for a node that is down After 2 minutes the copy is promoted; messages acked in the source's last moments may come back What the source committed after the destination's last read, if it never returns. On v3.0.1, records committed after the promote could also stay undelivered (fixed in v3.1.0) Keep any copy the returning source quarantines; see Scale out and in
A rolling restart Requests move to the other pods; Raft leadership moves in about 150 ms Messages leased on the restarting node come back; acked ones do not Nothing Roll one pod at a time
An unacked message outlives retention Nothing It is never delivered; the frontier skips it, and the owner logs it That message, by policy Alert on lag, and keep retention above your longest consumer outage
A remote child's target is unreachable (v3.2.0) Nothing Nothing here; the copy on the remote falls behind by the outage, then catches up, with some records twice Only records that age out of the parent before they ship, counted in narad_fanout_child_dropped_messages Alert on headroom, and keep the parent's retention above the longest outage you accept

Every status code, and whether to retry it, is in Status codes and errors.

How we check the contract

Every pull request runs a three-node cluster on loopback through two suites. One produces, consumes and acks a known set of messages with nothing broken, and fails on any message delivered twice or never delivered. The other does the same while nodes are stopped and restarted, and fails unless every message is acked. Both fail on a message the run never produced, or one on the wrong topic.

The property these runs check, the checks that run by hand, and what none of them can catch are in Linearizability.

Next steps