Blog

Queues and Logs: What Happens When the Consumer Dies

We killed Kafka consumers 310 times at random moments and counted every record lost or processed twice. The order of two lines decides which you get. Queues vs logs, delivery semantics, Kafka’s exactly-once, ordering, poison messages and backpressure.

A message broker sits between the service that has work and the service that does it. The producer writes a message and moves on. A consumer picks it up later, does the work and says it’s done. That gap is the point: the producer doesn’t wait, and a slow or broken consumer doesn’t stop it.

The gap is also where things go wrong. A consumer can die halfway through a batch. A producer can send a message, get no reply and send it again. This part is about what happens then. We ran Kafka 4.3.1 and RabbitMQ 4.3.6, killed Kafka consumers at random moments and RabbitMQ consumers mid-queue, and counted every message that was lost or processed twice.

Try this first

A consumer reads a batch of 500 messages. For each one, it charges a card. After the whole batch, it tells the broker “done up to here”. It’s killed at a random moment and restarted.

What happens?

  1. Nothing. The broker knows which messages were handled.
  2. Some cards are charged twice.
  3. Some cards are never charged.
  4. Either, depending on when the kill lands.

Write down your answer. The measured answer is below.

Queue or log

Message brokers come in two designs, queues and logs, and they make opposite bets.

A queue hands out a message, waits for the consumer to acknowledge it, then deletes it. RabbitMQ works this way. Its documentation says “all current RabbitMQ queue types have destructive consume behaviour, i.e. messages are deleted from the queue when a consumer is finished with them”. The broker tracks every message: delivered, acknowledged, or waiting to be delivered again.

A log keeps every message and never deletes on read. Kafka works this way. Its design document says the goal was something “more akin to a database log than a traditional messaging system”. Each reader keeps a bookmark, an offset, the position of the next message it wants. Kafka’s introduction says “events are not deleted after consumption”. They’re deleted by age instead, after seven days by default.

The difference shows up in what the broker has to remember. A queue remembers the state of each message. A log remembers one number per reader, per partition. Kafka’s design document puts it this way:

Our topic is divided into a set of totally ordered partitions, each of which is consumed by exactly one consumer within each subscribing consumer group at any given time. This means that the position of a consumer in each partition is just a single integer, the offset of the next message to consume.

Because the log keeps everything, a reader can go back. “A consumer can deliberately rewind back to an old offset and re-consume data,” the same document says. “This violates the common contract of a queue, but turns out to be an essential feature for many consumers.” A new service can read a year of history. A buggy consumer can be fixed and replayed.

The two are growing towards each other. RabbitMQ added streams, “an append-only log of messages that can be repeatedly read until they expire”. Kafka 4.3’s documentation describes share groups, where “Records are acknowledged individually” and “Delivery attempts to consumers in a share group are counted”. Pick by the shape of the work, not the product name. A queue suits independent jobs that anyone can take. A log suits a history that several readers each need all of.

Consumer groups and partitions

A Kafka topic is split into partitions. Each partition is its own log. A consumer group is a set of consumers that share the work. Inside a group, each partition goes to exactly one consumer. Across groups, every group gets everything.

That one rule gives both classic patterns. The Kafka Javadoc:

To get semantics similar to a queue in a traditional messaging system all processes would be part of a single consumer group and hence record delivery would be balanced over the group like with a queue. Unlike a traditional messaging system, though, you can have multiple such groups. To get semantics similar to pub-sub in a traditional messaging system each process would have its own consumer group, so each process would subscribe to all the records published to the topic.

It also sets a limit. A partition is read by one consumer per group, so a group can’t usefully have more consumers than the topic has partitions. The 2011 Kafka paper made this choice on purpose: “Our first decision is to make a partition within a topic the smallest unit of parallelism.” The proposal for share groups names the cost today:

Users of Kafka often have to “over-partition” simply to ensure they can have sufficient parallel consumption to cope with peak loads.

Delivery semantics: the order of two steps

A consumer does two things with each message. It does the work: charges the card, writes the row, sends the email. And it records its progress: commits an offset in Kafka, or acknowledges the message in RabbitMQ.

A crash can land between the two. Which one came first decides what the crash costs. Kafka’s design document:

It can read the messages, then save its position in the log, and finally process the messages. In this case there is a possibility that the consumer process crashes after saving its position but before saving the output of its message processing.

This corresponds to “at-most-once” semantics as in the case of a consumer failure messages may not be processed.

It can read the messages, process the messages, and finally save its position. In this case there is a possibility that the consumer process crashes after processing messages but before saving its position.

This corresponds to the “at-least-once” semantics in the case of consumer failure.

  • At-most-once: progress first, then work. A crash loses the messages in between. Nothing is done twice.
  • At-least-once: work first, then progress. A crash repeats the messages in between. Nothing is lost.
  • Exactly-once: each message’s effect happens once. Kafka defines it as “Each message is processed once and only once.” We’ll come back to what that takes.

Measured: killing the consumer

We ran one Kafka 4.3.1 broker in Docker and a consumer using the official Java client. The topic held 10,000 records in one partition. The consumer did 1 ms of work per record, then appended the record’s number to a file, one unbuffered write per record. That file stands in for a side effect you can’t take back. Once the write returns, the kernel has the data, so killing the process doesn’t undo it.

Then we killed the consumer with SIGKILL, at a random moment 0.5 to 6 seconds after its first record, and started it again. A killed process gets no chance to clean up: no final commit, no goodbye. The restarted consumer resumed from the group’s committed offset and ran to the end. It took its partition with assign() rather than joining the group with subscribe(), so a restart didn’t have to wait for the group to notice the dead member. Offsets are committed the same way either way. Then we counted the file: record numbers missing were lost, numbers written twice were duplicated. We did this 30 times for each way of committing, in shuffled order.

a consumer is killed mid-batch, then restarted records processed (side effect done) killed here committed offsetthe restart resumes here all 30 kills: 0 250 500 750 1,000 records lost or processed twice, per kill lost: never processedduplicated: processed twice

Measured by checks/part27_queues/kafka_lab.py: Kafka 4.3.1 and its Java client, one partition of 10,000 records, 1 ms of work per record, the consumer SIGKILLed 0.5 to 6 s after its first record. Lost and duplicated are counted in a file the consumer appends to, one unbuffered write per record. The consumer takes its partition with assign(), not subscribe(), so a restart doesn’t wait for the group. The worker thread’s queue holds 1,000 records.

So the answer to the opening question is 2. Committing after each batch of 500 duplicated records in 29 of 30 kills, up to 479 at once, and never lost one. Committing first lost records in all 30 kills, up to 484, and never duplicated one.

The size of the damage is the batch. The kill lands somewhere inside a batch, so the loss or duplication is anywhere from zero to one batch. With max.poll.records at its default of 500, the median was 283.5 lost or 239.5 duplicated. With batches of 50, it was 32.5 lost or 19 duplicated, and never more than 50.

Commit Batch Kills with loss Kills with duplicates Worst kill
Before the work 500 30 of 30 0 of 30 484 lost
Before the work 50 28 of 30 0 of 30 50 lost
After the batch 500 0 of 30 29 of 30 479 twice
After the batch 50 0 of 30 27 of 30 50 twice
After each record 0 of 30 1 of 30 1 twice

Committing after every record shrinks the window to one record. It isn’t free. With 1 ms of work per record, it ran at 423 records a second, against 914 for committing once per batch of 500. Each commit waits for a round trip to the broker. Even then, one kill in 30 landed between a write and its commit, and that record was processed twice. The window got smaller. It didn’t close.

Auto-commit

Kafka’s consumer commits for you by default. enable.auto.commit is true, and it commits every 5,000 ms. The Javadoc attaches a condition:

Note: Using automatic offset commits can also give you “at-least-once” delivery, but the requirement is that you must consume all data returned from each call to poll(Duration) before any subsequent calls, or before closing the consumer. If you fail to do either of these, it is possible for the committed offset to get ahead of the consumed position, which results in missing records.

We tested both sides of that condition.

Work finished before the next poll(). This is the plain loop: poll, handle every record, poll again. In 60 kills, at the default 5-second interval and at 100 ms, it never lost a record. It duplicated instead, and the interval wasn’t the bound.

  • At 100 ms, kills still repeated up to 416 records, about 0.4 seconds of work, and 25 of 30 repeated more than 100. None repeated more than one batch of 500. That fits a commit that only saves what earlier polls handed out, as the Javadoc’s condition suggests.
  • At the default 5 seconds, 21 of our 30 kills came before the first automatic commit, so the restart redid everything since the consumer started: a median of 2,886 records, and up to 4,497. That size comes from when we chose to kill, not from Kafka. The kills after a commit repeated 40 to 1,291 records.

Either way, a crash repeats everything since the last commit that actually happened, and that can be more than one interval.

Work handed to another thread. This is a common way to speed up a slow consumer: the polling thread puts records on a queue, and a worker takes them off. poll() returns quickly and is called again, and auto-commit saves the position of what poll() handed out, not what the worker finished. At a 100 ms interval, all 30 kills lost records, 515 to 985 each. Our worker’s queue held 1,000 records, and every loss fell between 500 and 1,000, never more than the queue held. That fits a commit that saved the position of records earlier polls had already queued, which the worker hadn’t reached yet.

At the default 5-second interval, the same design lost records in 6 kills and duplicated in the other 24. The 6 losses came from kills 3.7 to 5.2 seconds after the first record, around each trial’s first automatic commit, whose timing varied from trial to trial. We didn’t log commit times. Right after a commit, the saved position is ahead of the worker, and a kill loses the gap. Kills before any commit redid everything, and the last few, after the worker had caught up, redid a little. The split of 6 and 24 comes from when we chose to kill, not from a rate. Which one you get depends on when the crash lands. You don’t get to choose.

The fix is what the Javadoc says: “Typically, you must disable automatic commits and manually commit processed offsets for records only after the thread has finished handling them”.

The producer side: a timeout isn’t a no

The consumer isn’t the only source of duplicates. A producer that sends a message and gets no reply can’t tell whether the broker wrote it. Kafka’s design document:

If a producer attempts to publish a message and experiences a network error, it cannot be sure if this error happened before or after the message was committed.

If it sends again and the first one had landed, the log has the message twice. This is the same problem as a timed-out write in Part 23.

Since version 0.11, Kafka’s producer can remove these duplicates. The broker gives each producer an ID, the producer numbers what it sends, and the broker drops what it has already written. Kafka’s Javadoc says this idempotent producer has been on by default since Kafka 3.0.

We measured it. We added a random delay to everything the broker sent, with tc netem: a 20 ms delay with 60 ms of jitter. The producer’s request timeout was 100 ms. We added no loss, and the delay applied only to what the broker sent, so requests arrived on time; the answers were late. We ran it 16 times, 8 with idempotence off and 8 with it on. Each run sent the numbers 0 to 19,999 and then read the partition back.

  • Idempotence off: every run had duplicates, 5,345 to 21,751 of them on 20,000 sends. In all 8 runs, the number of duplicates equalled the producer’s own count of retried records exactly. Every retry was a record the broker had already written.
  • Idempotence on: the producer retried 1,371 to 12,766 records per run, and the log held each number exactly once, in all 8 runs.

So in the idempotence-off runs, every record the producer retried had in fact been written: the worst case for duplicates. The lab’s point is the second result: the idempotent producer removed all of them. How many retries you get depends on your network. A request that’s lost before the broker reads it retries without duplicating, and this lab didn’t create any.

Kafka’s documentation also warns that retries with idempotence off and several requests in flight “will potentially change the ordering of records”. That needs an earlier batch that failed without being written, which this lab never produced, so it says nothing either way. The documented rule is the one to rely on: “if retries are disabled or if enable.idempotence is set to true, ordering will be preserved.”

Two limits are worth knowing. Idempotence covers the producer’s own retries, “within a single session”. If your code calls send() twice for the same thing, that’s two records, and the Javadoc says “it is imperative to avoid application level re-sends since these cannot be de-duplicated.” And a conflicting setting can switch it off quietly. From the producer configuration: “If conflicting configurations are set and idempotence is not explicitly enabled, idempotence is disabled.”

Exactly-once, and what it covers

“Exactly-once delivery is impossible” and “Kafka has exactly-once” are both said often. They’re talking about different things.

The impossibility argument is about acknowledgements. Kyle Kingsbury put it plainly in his 2014 Jepsen analysis of RabbitMQ:

Acknowledge before processing the message, and a crash can cause data loss. Acknowledge after processing, and a crash can cause duplicate delivery. No distributed queue can offer exactly-once delivery–the best they can do is at-least-once or at-most-once.

That’s what our first figure measured. Tyler Treat’s 2015 post “You Cannot Have Exactly-Once Delivery” also says how systems get around it: “The way we achieve exactly-once delivery in practice is by faking it.” Messages are made safe to apply twice, or duplicates are filtered out.

What Kafka offers is narrower and real. Its design document:

As a result, Kafka supports exactly-once delivery in Kafka Streams, and the transactional producer and the consumer using read-committed isolation level can be used generally to provide exactly-once delivery when reading, processing and writing data on Kafka topics. Exactly-once delivery for other destination systems generally requires cooperation with such systems, but Kafka provides the primitives which makes implementing this feasible (see also Kafka Connect). Otherwise, Kafka guarantees at-least-once delivery by default, and allows the user to implement at-most-once delivery by disabling retries on the producer and committing offsets in the consumer prior to processing a batch of messages.

The trick is that Kafka stores consumer offsets in a Kafka topic too. So a program that reads from one topic and writes to another can put its output records and its input offset in one transaction. Either both commit or neither does. If it crashes mid-batch, that transaction never commits. The restarted program should use the same transactional.id, because that’s how Kafka finds the dead one’s work at once. The Javadoc for initTransactions(): “If the previous instance had failed with a transaction in progress, it will be aborted.” If nothing comes back, the transaction coordinator aborts it anyway, within transaction.timeout.ms of the transaction starting, 60 seconds by default. Whoever takes the partition next redoes the batch, and readers that ask for committed data never see the unfinished copy. Until the abort, those readers wait: they only read up to the first open transaction.

Readers who ask, anyway. The consumer default is read_uncommitted, and “consumer.poll() will return all messages, even transactional messages which have been aborted.” Only readers set to read_committed get the guarantee.

We measured all three places a copy could show up. The worker read 5,000 records, appended each to a file (the side effect outside Kafka), and wrote each to an output topic. We killed it once per trial at a random moment and restarted it, the transactional runs with the same transactional.id, 20 times with a transaction and 20 without.

consume, do a side effect, produce to another topic; killed once records there twice (median of 20 kills) a file (outside Kafka) topic, read_uncommitted topic, read_committed 0 250 500

Measured by checks/part27_queues/kafka_lab.py: Kafka 4.3.1, 5,000 input records, batches of 500, the worker SIGKILLed 0.3 to 4 s after its first record, then restarted; the transactional runs restart with the same transactional.id. The side effect is an append to a file.

With the transaction, a read_committed reader saw every record exactly once in all 20 trials. A reader left at the default saw duplicates in all 20, a median of 291. And the file had duplicates in all 20, a median of 293.5. The transaction was aborted, and Kafka hid its records from readers who asked it to. It couldn’t un-append a file.

Confluent’s post on exactly-once, bylined Neha Narkhede, Guozhang Wang and Confluent Staff, first published in 2017 and updated in March 2025, says the same about Kafka Streams:

Note that exactly-once semantics is guaranteed within the scope of Kafka Streams’ internal processing only; for example, if the event streaming app written in Streams makes an RPC call to update some remote stores, or if it uses a customized client to directly read or write to a Kafka topic, the resulting side effects would not be guaranteed exactly once.

The mechanism has had bugs of its own. Jepsen’s 2024 analysis of Bufstream, a Kafka-compatible system, reported: “We observed aborted reads and torn transactions due to process pauses in Kafka as well, and opened KAFKA-17754 to track the issue.” The issue is marked fixed, but its record names no release that carries the fix.

So Kafka’s exactly-once means this: a read-process-write loop whose input and output both live in Kafka commits atomically. Anything it does to the outside world is at-least-once. Charging a card, sending an email or writing to another database needs the same thing it always needed: an idempotency key, or the offset stored in the same transaction as the output. Kafka’s design document suggests the second: “letting the consumer store its offset in the same place as its output.”

The other brokers draw the same line in their own words. Amazon SQS FIFO queues offer “exactly-once processing”, which AWS defines this way: “If you retry the SendMessage action within the 5-minute deduplication interval, Amazon SQS doesn’t introduce any duplicates into the queue.” The same guide warns: “If a consumer fails to process a message before the visibility timeout expires, another consumer may receive and process it, leading to duplicate processing.” RabbitMQ’s streams documentation says that deduplicated publishing, using the source message’s offset as the publishing ID, makes a read-process-write loop “effectively-once (what other systems call exactly-once)”. Its queues make no such promise.

Ordering: per partition, and only by key

Kafka keeps order within a partition. Its introduction: “Kafka guarantees that any consumer of a given topic-partition will always read that partition’s events in exactly the same order as they were written.” Across partitions, there’s no order. The 2011 paper says so directly: “there is no guarantee on the ordering of messages coming from different partitions.”

The way to keep related messages in order is the key. “Events with the same event key (e.g., a customer or vehicle ID) are written to the same partition”. Give every event for one account the same key, and they share a partition, and so an order.

We wrote 20 accounts × 500 numbered events into a topic with 6 partitions, interleaved, and read them back with a group of 1 consumer and a group of 3. Every consumer did 0.2 ms of work per event and appended it to one shared file, standing in for a shared database downstream. Then we counted events that landed in the file after a later event for the same account. Five runs of each. Each consumer recorded which partitions it read. With keys, every one of the 3 took 2 partitions and handled 3,000 or 3,500 events. Without keys, in one run one of the 3 consumers got nothing at all, because the events had landed in only 3 of the 6 partitions.

Events Consumers Out-of-order events (of 10,000) Accounts affected (of 20)
Keyed by account 1 0 in every run 0
Keyed by account 3 0 in every run 0
No key 1 2,701 to 8,248 20 in every run
No key 3 6,135 to 8,478 20 in every run

Without a key, the producer spread each account’s events over 3 to 6 of the 6 partitions, and every account came out of order in every run. With a key, not one event in 100,000 did.

Pick the key by what must stay in order. Order within an account needs the account ID. A global order needs one partition, which means one consumer at a time. The same shape shows up elsewhere: SQS FIFO orders “within each message group”, and Google Pub/Sub only orders messages that share an ordering key.

Poison messages and dead-letter queues

Some messages can never be processed. The payload is malformed, or it triggers a bug, or it refers to something deleted. Each time it’s delivered, the consumer fails. What the broker does next decides whether one bad message stalls the whole queue.

RabbitMQ’s quorum queues count failed deliveries. Their documentation:

Quorum queues keep track of the number of unsuccessful (re)delivery attempts and expose it in the “x-delivery-count” header that is included with any redelivered message.

When a message has been redelivered more times than the limit the message will be dropped (removed) or dead-lettered (if a DLX is configured).

“Starting with RabbitMQ 4.0, the delivery limit for quorum queues defaults to 20.” A dead-letter exchange (DLX) is where RabbitMQ republishes a message it gives up on, so it lands in a dead-letter queue someone can look at. That republish isn’t guaranteed either: “at-most-once remains the default dead-letter-strategy for quorum queues”, so in a cluster a dead-lettered message can be lost on the way unless you switch to the at-least-once strategy, which has requirements of its own.

We put one poison message at the front of a queue, followed by 100 good ones, and gave the consumer three ways to handle the poison: crash, reject it, or nack it with requeue=true, which means “not now, give it back”. A crash was the consumer process exiting without an acknowledgement, so RabbitMQ saw its connection drop. We tried a quorum queue with a DLX, a quorum queue without one, and a classic queue, at prefetch 1 and prefetch 10. Each case ran three times, and the three runs agreed on where every message ended up.

one message the consumer can never handle, then 100 good ones queue (front on the right) consumer processed (of 100 good) dead-letter queue deleted times the poison was delivered: x-delivery-count on the last one:

Measured by checks/part27_queues/rabbit_lab.py: RabbitMQ 4.3.6 with its default delivery limit, the pika client, each consumer its own process; a crash is the process exiting without an ack. Shown: the poison and the first 10 of the 100 good messages; three runs of each case, which agreed on where every message ended up.

What we saw on RabbitMQ 4.3.6:

  • Crash, quorum queue. The poison was delivered 21 times. The redeliveries carried x-delivery-count 1 to 20, and then it was dead-lettered with the reason delivery_limit. At prefetch 1, the good messages waited behind it, then all 100 were processed. Prefetch 10 is below.
  • Reject, quorum queue. One delivery, straight to the dead-letter queue with the reason rejected. This is the cheapest way out, if the consumer can tell the message is bad.
  • Nack with requeue, quorum queue. The count never moved. The poison was delivered 828 to 1,254 times in 5 seconds, with no x-delivery-count at all, and it was still there when we stopped. The documentation says so: requeues through nack “will not count toward the delivery limit, allowing unlimited requeue loops for application-level routing logic”. Its table of what counts lists basic.reject as counting, so reject with requeue would reach the limit where nack never does. We didn’t test that one. At prefetch 1, not one good message got past the looping poison. At prefetch 10, all 100 did, around it, while the poison went on looping: 693 to 925 deliveries in 5 seconds.
  • Classic queue. No delivery limit. The poison came back after every crash, and we stopped at 40. At either prefetch, no good message was processed.

Two results are easy to miss.

The poison took its neighbours with it. At prefetch 10, the consumer held 10 messages when it crashed: the poison and the next 9. A crash counts against every message the consumer was holding, so all 10 reached the limit together. The 9 good ones were dead-lettered with the poison, never processed, with the same reason, delivery_limit. The documentation does warn: “For consumers with a prefetch greater than 1, all outstanding unacknowledged deliveries can be discarded by this mechanism if they all were requeued as a group multiple times, due to, say, consumer application instance failures.”

Without a DLX, the limit means deletion. The same case with no dead-letter exchange deleted the poison and 9 good messages. 91 of 100 good messages were processed. Every publish had already been confirmed, and nothing afterwards told any client about the other 9. They were just gone. The documentation warns: “It is recommended that all quorum queues have a dead letter configuration of some sort to ensure messages aren’t dropped and lost unintentionally.”

Other brokers have the same tool with different defaults, and none of them turns it on for you. Give an SQS queue a dead-letter queue and it moves a message there after maxReceiveCount receives, 10 by default. Give a Google Pub/Sub subscription a dead-letter topic and it defaults to 5 attempts, a maximum it says “is approximate”. Kafka’s ordinary consumer groups have none, since a log doesn’t track single messages. A consumer that can’t skip a bad record stops at it, or skips it on purpose and writes it to a topic of its own.

A dead-letter queue is a quarantine. It stops one bad message blocking the rest. It doesn’t fix anything, and someone has to read it.

Backpressure

A producer can outrun a consumer for as long as it likes. Something has to give. Jay Kreps listed the options in his 2013 essay on logs: processing “must block, buffer or drop data.”

A log buffers. Kafka consumers pull, so a slow one doesn’t get flooded. Kafka’s design document:

A pull-based system has the nicer property that the consumer simply falls behind and catches up when it can.

The messages wait in the log, on disk, for as long as retention allows: seven days by default. A consumer that falls further behind than that loses messages to retention instead. Watch consumer lag, the gap between the newest offset and the group’s committed offset, and alert on it.

A queue pushes. RabbitMQ sends messages to the consumer, and prefetch is the limit on how many can be delivered but not yet acknowledged. The documentation:

A value of 0 means “no limit”, allowing any number of unacknowledged messages.

The server’s default is 0, unless a client library sets its own.

We measured it. We published 5,000 messages to a quorum queue, and the consumer took 2 ms per message. Two seconds in, we killed the consumer and counted what RabbitMQ put back in the queue: every message it had delivered and not yet had acknowledged. Then a new consumer drained the queue, and we counted messages the application saw twice. Separately, we timed how fast the queue drained with no kill. Three runs per setting.

Prefetch Handed back by the kill, flagged redelivered Seen twice by the app Drain rate (msgs/s)
1 1 0 to 1 125 to 145
10 10 1 342 to 377
100 100 0 to 1 284 to 368
1,000 1,000 0 to 1 332 to 384
0 (no limit) 2,000 0 to 1 332 to 379

Three things stand out.

  • “No limit” stopped at 2,000. The documentation explains: “quorum queues cap consumer prefetch maximum to 2,000 to limit a potentially runaway Raft log growth”. We only measured quorum queues, so this is the quorum queue’s limit, not a general one.
  • Prefetch 1 was slow. It drained at 125 to 145 messages a second, against 284 to 384 for 10 and up. With 2 ms of work, the ceiling is 500. The consumer can’t get its next message until the broker has handled the last ack, and on this quorum queue a message took 6.9 to 8.0 ms at prefetch 1, against 2.7 to 2.9 ms at prefetch 10; we didn’t measure where the difference went. The documentation warns that prefetch 1 “will significantly reduce throughput”, and says values from 100 to 300 “usually offer optimal throughput and do not run significant risk of overwhelming consumers.”
  • The redelivered flag is only a hint. After each kill, RabbitMQ requeued exactly what the consumer had held, and flagged all of it redelivered: 2,000 messages at no limit. The consumer’s handler had run on at most 1 of them. The documentation says so: “This is a hint that a consumer may have seen this message before.” Don’t use it to skip work.

Prefetch is backpressure for consumers. For producers, RabbitMQ has its own: it “will reduce the speed of connections which are publishing too quickly for queues to keep up”, and when memory or disk runs low, it blocks publishers. Kafka’s producer blocks send() when its 32 MiB buffer fills, then fails after max.block.ms, 60 seconds by default.

The rule under all of these is stated well by the Reactive Streams specification: “all buffer sizes are to be bounded and these bounds must be known and controlled by the subscribers.” A buffer with no limit doesn’t remove the problem. It turns it into a memory problem somewhere else.

Explain it like I’m ten

A class has a box of homework to mark, and several helpers.

  • A queue is a box you take papers out of. A helper takes a paper, marks it and throws it away. If a helper faints holding a paper, it goes back in the box.
  • A log is a notebook nobody tears pages out of. Each helper keeps a bookmark. Anyone can go back and read an old page again.
  • Order matters for bookmarks. If you move your bookmark before marking the page, and you faint, that page never gets marked. If you mark first and faint before moving the bookmark, you mark it again when you wake up.
  • A paper nobody can mark gets put aside in a special pile after too many tries, so it doesn’t stop everyone.
  • Each helper only takes a few papers at a time, so nobody drowns in them.

The precise version

  • The box is a queue with destructive consume. The notebook is a log, and the bookmark is a committed offset.
  • Bookmark first is at-most-once. Mark first is at-least-once. Measured with batches of 500: the first lost records in every kill, the second duplicated them in 29 of 30.
  • Exactly-once is possible only inside a system that owns both the work and the bookmark. Kafka’s transactions give it for Kafka-to-Kafka processing, read with read_committed. Side effects outside need idempotency.
  • The special pile is a dead-letter queue, and “too many tries” is the delivery limit.
  • Taking a few papers at a time is prefetch, or more generally backpressure.
  • Where the analogy breaks: a fainting helper in RabbitMQ counts as a failed try for every paper they were holding, not just the one they were reading. Measured: 9 good messages were dead-lettered with one bad one.

Trade-offs

  • At-most-once or at-least-once is a choice between two failures. Most systems pick at-least-once and make the work idempotent.
  • Smaller batches mean less damage per crash and more commits. Measured: committing every record halved throughput and still duplicated one record in 30 kills.
  • A log replays; a queue forgets. Keeping history costs disk and gives you reprocessing. Deleting on ack keeps the broker small.
  • Partitions buy parallelism and cost global order. Measured: with keys, no event in 100,000 was out of order; without, every account was.
  • Transactions give exactly-once inside Kafka, and only for readers set to read_committed. Everything outside still needs idempotency.
  • Prefetch trades throughput against what a crash hands back. Measured on a quorum queue: prefetch 1 drained at under half the rate of prefetch 10.

Common mistakes

  • Handing records to a thread pool with auto-commit on. Measured: every kill at a 100 ms interval lost 515 to 985 records.
  • Believing auto-commit is at-most-once, or at-least-once, whatever you do. It’s at-least-once only if the batch is finished before the next poll(), and a crash can repeat more than one interval’s work.
  • Assuming “exactly-once” covers your database or email. Measured: with a Kafka transaction, the side effect outside Kafka still ran twice for a median of 293.5 records per kill.
  • Reading a transactional topic with the default read_uncommitted. Measured: aborted records showed up in every trial.
  • Retrying a send in application code. The idempotent producer can’t de-duplicate your own re-sends. Use an idempotency key the consumer checks.
  • Expecting order across partitions. Pick the key by what must stay in order.
  • Using nack with requeue for a message that will never succeed. Measured: 693 to 1,254 deliveries in 5 seconds on a quorum queue, and the limit never reached.
  • Running quorum queues without a dead-letter exchange. Measured: 9 good messages deleted along with one bad one.
  • Setting prefetch to 0 and calling it unlimited. Measured: on a quorum queue it quietly meant 2,000, and a crash handed all 2,000 back.
  • Trusting the redelivered flag. Measured: 2,000 messages flagged, at most 1 actually seen twice.

Interview questions

1. What’s the difference between a message queue and a log?

A queue deletes a message once a consumer acknowledges it, and tracks the state of each message. A log keeps every message for a retention period and tracks one offset per consumer group per partition. A log can replay history to a new or fixed consumer; a queue can’t. RabbitMQ queues are the first kind; Kafka topics, and RabbitMQ streams, are the second.

2. Explain at-most-once, at-least-once and exactly-once.

At-most-once records progress before doing the work, so a crash can lose messages. At-least-once does the work first, so a crash can repeat it. Exactly-once means each message’s effect happens once. Over a network where acknowledgements can be lost, you get it by making repeats harmless, with idempotent work or deduplication, or by committing the output and the progress in one transaction.

3. Does Kafka provide exactly-once delivery?

For a read-process-write loop whose input and output are both Kafka topics, yes: the transactional producer commits the output records and the input offsets atomically, and consumers must read with read_committed. For side effects outside Kafka, no. The default is at-least-once. Measured: with transactions, read_committed readers saw no duplicates in 20 kills, and a file the same code appended to had duplicates in all 20.

4. A consumer uses auto-commit. Can it lose messages?

Yes, if it processes records outside the poll loop, for example on a worker thread. Auto-commit saves the position of what poll() returned, which can be ahead of what was processed. If every record is finished before the next poll(), it’s at-least-once, and a crash repeats everything since the last commit that actually happened. Measured with a 100 ms interval: repeats of up to 416 records, about 0.4 seconds of work.

5. How do you keep messages in order in Kafka?

Order is guaranteed only within a partition. Give related messages the same key so they land in the same partition, and have one consumer per partition, which a consumer group does for you. Use the idempotent producer, the default since Kafka 3.0, so retries don’t reorder. The producer picks the partition from “a hash of the key”, so don’t change the partition count casually: with a different count, a key can map to a different partition. A global order needs a single partition.

6. How do you handle a poison message?

Put a limit on delivery attempts and a dead-letter queue behind it. Reject a message the consumer can tell is bad, rather than requeue it, and alert on the dead-letter queue. Watch what a crash counts: on a RabbitMQ quorum queue, a consumer crash counts against every message it held, so a low prefetch limits the collateral.

7. What is backpressure, and how do Kafka and RabbitMQ apply it?

Backpressure is a slow consumer limiting how fast work arrives, so buffers stay bounded. Kafka consumers pull, so a slow one falls behind and the backlog waits in the log; you watch consumer lag. RabbitMQ pushes, bounded by prefetch on the consumer side, and slows or blocks publishers when it can’t keep up. Producers block when their own buffer fills.

Sources

What to remember

  • A queue deletes on acknowledgement. A log keeps everything and gives each reader a bookmark.
  • The order of “do the work” and “record progress” decides the failure. Measured: commit first lost records in 30 of 30 kills; commit after duplicated them in 29 of 30. The damage is up to one batch.
  • Auto-commit is at-least-once only if you finish each batch before the next poll(). Measured: with a worker thread, every kill at 100 ms lost records.
  • Kafka’s exactly-once covers Kafka-to-Kafka processing read with read_committed. Measured: no duplicates there in 20 kills, and duplicates outside Kafka in all 20.
  • The idempotent producer removes duplicates from retries. Measured: under a delay that made every retry a duplicate with idempotence off, idempotence left none in 8 runs.
  • Order lives inside a partition, so choose the key by what must stay in order.
  • A dead-letter queue is a quarantine, not a fix. Measured: on a RabbitMQ quorum queue, a crash dead-lettered 9 good messages with the bad one, and without a DLX it deleted them.
  • Bound every buffer. Prefetch 0 isn’t unlimited on a quorum queue; it’s 2,000.

Every broker is at-least-once somewhere. Design the work so that doing it twice is harmless, and most of this part stops mattering.

How useful was this post?

Click on a heart to rate it!

Average rating 0 / 5. Vote count: 0

No votes so far! Be the first to rate this post.