Learn Labs
7. Reliable Data Delivery

7.4 Using producers reliably

That parenthetical is worth quoting to anyone waving a Kafka benchmark at you.

"Even if we configure the brokers in the most reliable configuration possible, the system as a whole can still potentially lose data if we don't configure the producers to be reliable as well."

4.1 The two loss scenarios — with the perfect broker configuration

Scenario 1 — acks=1 loses an acknowledged message

Broker config: 3 replicas, unclean leader election disabled — “So we should never lose a single message that was committed to the Kafka cluster.” Producer config: acks=1 ◄ the flaw.

replica.lag.time.max.msmessage written to the LEADERnot yet to the in-sync replicas“Message was written successfully”the leader’s response to the producerleader IMMEDIATELY CRASHESbefore the data was replicatedthe other replicas are STILL IN SYNCit takes a while before we declare a replica out of syncone of them becomes leaderthe message was never written to itMESSAGE LOSTthe producing application thinks it was written

Note the precise accounting: “The system is CONSISTENT because no consumer saw the message (it was never committed because the replicas never got it), but FROM THE PRODUCER PERSPECTIVE, A MESSAGE WAS LOST.”

Scenario 2 — acks=all plus bad error handling still loses

Broker config: 3 replicas, unclean leader election disabled. Producer config: acks=all — correct!

we attempt to writethe leader JUST CRASHED; a new one is being electedKafka responds “Leader not Available”the producer MISHANDLES the errorand doesn’t retryMESSAGE MAY BE LOST

“this is NOT a broker reliability issue because the broker never got the message; and it is NOT a consistency issue because the consumers never got the message either. But if producers don’t handle errors correctly, they may cause message loss.”

Figure 7.4.14.1 The two loss scenarios — with the perfect broker configuration

The two rules that follow:

 ① Use the correct `acks` configuration to match reliability requirements
 ② HANDLE ERRORS CORRECTLY — both in configuration AND in code

4.2 The three ack modes, restated with reliability precision

What it meansWhat you lose
acks=0"a message is considered written successfully if the producer managed to send it over the network""We will still get errors if the object cannot be serialized or if the network card failed, but we won't get any error if the partition is offline, a leader election is in progress, or even if the ENTIRE KAFKA CLUSTER IS UNAVAILABLE."
acks=1"the leader will send either an acknowledgment or an error the moment it gets the message and writes it to the partition data file (but not necessarily synced to disk)"Scenario 1 above. Plus: "it is also possible to write to the leader faster than it can replicate messages and end up with under-replicated partitions, since the leader will acknowledge messages from the producer before replicating them."
acks=all"the leader will wait until all in-sync replicas get the message before sending back an acknowledgment or an error. In conjunction with min.insync.replicas, this lets us control how many replicas get the message before it is acknowledged.""the option with the longest producer latency"

💡 The acks=0 benchmark warning

"Running with acks=0 has low produce latency (which is why we see a lot of benchmarks with this configuration), but it will NOT improve end-to-end latency (remember that consumers will not see messages until they are replicated to all available replicas)."

That parenthetical is worth quoting to anyone waving a Kafka benchmark at you. And it restates Ch. 3's core insight: lower acks buys producer-side latency only, never end-to-end latency.

4.3 Producer retries

Retriable vs non-retriable, with concrete error codes:

ClassExampleWhy
RetriableLEADER_NOT_AVAILABLE“maybe a new broker was elected and the second attempt will succeed”
Non-retriableINVALID_CONFIG“trying the same message again will not change the configuration”

The recommended configuration:

"when our goal is to never lose a message, our best approach is to configure the producer to keep trying to send the messages when it encounters a retriable error. And the best approach to retries... is to leave the number of retries at its current default (MAX_INT, or effectively infinite) and use delivery.timeout.ms to configure the maximum amount of time we are willing to wait until giving up — the producer will retry sending the message as many times as possible within this time interval."

⚠️ Retries guarantee at-least-once, NOT exactly-once

"Retrying to send a failed message includes a risk that both messages were successfully written to the broker, leading to duplicates. Retries and careful error handling can guarantee that each message will be stored AT LEAST ONCE, but not EXACTLY ONCE. Using enable.idempotence=true will cause the producer to include additional information in its records, which brokers will use to skip duplicate messages caused by retries."

Reliable producer =

  • acks=all
  • retries left at default (~infinite)
  • delivery.timeout.ms greater than cluster recovery time
  • enable.idempotence=true
  • min.insync.replicas=2 on the broker side

Without idempotence you have at-least-once. That’s a correct choice — but you must know which one you chose. (→ Ch. 8)

4.4 Errors YOU must handle

"Using the built-in producer retries is an easy way to correctly handle a large variety of errors without loss of messages, but as developers, we must still be able to handle other types of errors:"

  • Non-retriable broker errors — message size, authorization, etc.
  • Errors before the message was sent to the broker — e.g. serialization
  • Errors when the producer exhausted all retry attempts, or when available memory is filled to the limit “due to using all of it to store messages while retrying”
  • Timeouts

The design questions the book poses — note it refuses to answer them for you:

"do we throw away 'bad messages'? Log errors? Stop reading messages from the source system? Apply back pressure to the source system to stop sending messages for a while? Store these messages in a directory on the local disk? These decisions depend on the architecture and the product requirements."

The one clear rule: "Just note that if all the error handler is doing is retrying to send the message, then we'll be better off relying on the producer's retry functionality."

Hand-rolled retry loops on top of the producer's retry machinery are strictly worse: they duplicate delivery, break the delivery.timeout.ms accounting, and can reorder messages.


On this page