Learn Labs
7. Reliable Data Delivery

7.3 Broker configuration — three knobs

replication.factor (topic) / default.replication.factor (broker, for auto-created topics).

Durability settings

acks
unclean election
leaderin ISR
replica 1in ISR
replica 2down
2/3
in ISR
2
required
Safe

Durable and available. 2 in-sync replicas meets the minimum of 2, and acks=all waits for all of them.

Durability is the product of three settings, not any one of them: acks on the producer, replication factor on the topic, and min.insync.replicas on the broker. Changing one in isolation is how clusters silently lose data.

"these can apply at the broker level, controlling configuration for all topics, and at the topic level, controlling behavior for a specific topic."

Why per-topic control matters — the bank example:

"at a bank, the administrator will probably want to set very reliable defaults for the entire cluster but make an exception to the topic that stores customer complaints where some data loss is acceptable."

3.1 Replication factor

replication.factor (topic) / default.replication.factor (broker, for auto-created topics).

A replication factor of N:

  • ✓ allows us to lose N−1 brokers while still reading and writing
  • ✗ requires at least N brokers
  • ✗ stores N copies → N times as much disk space

“We are basically trading availability for hardware.”

"Even after a topic exists, we can choose to add or remove replicas and thereby modify the replication factor using Kafka's replica assignment tool."

The five considerations for choosing N
ConsiderationThe argument
Availability"A partition with just one replica will become unavailable even during a ROUTINE RESTART of a single broker."
Durability"If a partition has a single replica and the disk becomes unusable... we've lost all the data. With more copies, especially on different storage devices, the probability of losing all of them is reduced."
ThroughputReplication traffic multiplies — see the table below
End-to-end latency"Each produced record has to be replicated to all in-sync replicas before it is available for consumers. In theory, with more replicas, there is higher probability that one of these replicas is a bit slow... In practice, if one broker becomes slow for any reason, it will slow down every client that tries using it, REGARDLESS OF REPLICATION FACTOR."
Cost"the most common reason for using a replication factor lower than 3 for noncritical data"

The throughput arithmetic — worth internalizing:

produce rate: 10 MBps to a partition
1 replica
0 MBps replication
2 replicas
10 MBps replication
3 replicas
20 MBps replication
5 replicas
40 MBps replication
  • · replication traffic = produce_rate × (RF − 1)
  • · “We need to take this into account when planning cluster size and capacity.”
Figure 7.3.2The throughput arithmetic — worth internalizing

(This is the term Ch. 2 said people forget when sizing NICs: replication is an additional consumer of every byte.)

The RF=2 cost argument, and its honest caveat:

"Since many storage systems already replicate each block 3 times, it sometimes makes sense to reduce costs by configuring Kafka with a replication factor of 2. Note that this will still REDUCE AVAILABILITY compared to a replication factor of 3, but DURABILITY will be guaranteed by the storage device."

That's a precise distinction: underlying storage replication protects your bytes; it does nothing for partition availability during a broker restart or failure.

Replica placement

"Kafka will always make sure each replica for a partition is on a separate broker. In some cases, this is not safe enough. If all replicas for a partition are placed on brokers that are on the same rack, and the top-of-rack switch misbehaves, we will lose availability of the partition REGARDLESS OF THE REPLICATION FACTOR."

P0 leaderP0 followerP0 followerone bad ToR switchpartition OFFLINERACK A — RF = 3, all three replicas here

RF=3 bought you nothing against this failure. The fix is broker.rack on every broker: “If rack names are configured, Kafka will make sure replicas for a partition are spread across multiple racks.” In cloud, “it is common to consider availability zones as separate racks.”

Figure 7.3.3Replica placement

(Mechanism: Ch. 6 §6.3's rack-alternating broker list. Caveat: Ch. 2 — rack awareness applies to newly created partitions only, and nothing monitors it after a reassignment.)

3.2 Unclean leader election

unclean.leader.election.enable — broker-level only (and in practice cluster-wide). Default: false.

"leader election is 'clean' in the sense that it guarantees no loss of committed data — by definition, committed data exists on all in-sync replicas. But what do we do when no in-sync replica exists except for the leader that just became unavailable?"

The two scenarios that produce this state
Scenario A — followers crash, then the leader crashes
the one and only in-sync replicaan out-of-sync follower starts firstpartition has 3 replicasthe TWO FOLLOWERS become unavailabletwo brokers crashproducers keep writing to the leaderall messages ACKNOWLEDGED AND COMMITTEDthe LEADER becomes unavailableanother crashan out-of-sync replicais the ONLY available replica

acks=all provided no protection here — “since the leader is the one and only in-sync replica.”

Scenario B — network degradation, then the leader crashes
network issuespartition has 3 replicasthe two followers FALL BEHIND“up and replicating, but NO LONGER IN SYNC”the leader keeps accepting messagesas the only in-sync replicathe leader becomes unavailableonly out-of-sync replicas remain
Figure 7.3.4The two scenarios that produce this state

Scenario A is the one that should change how you configure Kafka. With min.insync.replicas=1 (the default), acks=all degenerates into acks=1 the moment your followers die — and it does so silently. This is precisely the hole min.insync.replicas=2 plugs (§3.3).

The choice, and the consistency damage spelled out
Don’t allow out-of-sync leaders (default)Do allow it
“the partition will remain offline until we bring the old leader (and the last in-sync replica) back online. In some cases (e.g., memory chip needs replacement), this can take many hours.”“we are going to lose all messages that were written to the old leader while that replica was out of sync and also cause some inconsistencies in consumers.”

The book's worked example of the inconsistency — this is the part people underestimate:

While replicas 0 and 1 were unavailable, offsets 100–200 were written to replica 2, the leader.

unclean leader electionreplica 2 — the leaderoffsets 100–200 written herereplica 2 unavailablereplica 0 comes back onlinehas 0–100 but NOT 100–200replica 0 becomes leaderproducers write COMPLETELY NEW 100–200replica 2 returns as a FOLLOWERDELETES messages not on the current leader
  • · Consumer fallout: “SOME consumers may have read the OLD messages 100–200, SOME consumers got the NEW 100–200, and SOME GOT A MIX OF BOTH.” “This can lead to pretty bad consequences when looking at things like downstream reports.”
  • · And then replica 2 returns: “it will become a FOLLOWER of the new leader. At that point, it will DELETE any messages it got that don’t exist on the current leader. Those messages will NOT BE AVAILABLE TO ANY CONSUMER IN THE FUTURE.”
Figure 7.3.6The book's worked example of the inconsistency — this is the part people underestimate
replica 2… 98 99OLD dataacknowledged, then DELETEDreplica 0… 98 99NEW datadifferent bytes, same offsetsnewoffsets … 98 99offsets 100 … 200offsets 201 …

The same offset now means two different messages depending on which consumer read when. Offsets stopped being a stable identity.

Figure 7.3.7

The operational recipe:

"By default it is set to false, which will not allow out-of-sync replicas to become leaders. This is the safest option... It is always possible for an administrator to look at the situation, decide to accept the data loss in order to make the partitions available, and switch this configuration to true before starting the cluster. JUST DON'T FORGET TO TURN IT BACK TO false AFTER THE CLUSTER RECOVERED."

(Ch. 5's electLeaders(ElectionType.UNCLEAN) is the per-partition, no-restart alternative.)

3.3 min.insync.replicas — the fix for the silent-degradation hole

Available at both topic and broker level.

The precise statement of the problem:

"part of the problem is that, per Kafka reliability guarantees, data is considered committed when it is written to all in-sync replicas — EVEN WHEN 'ALL' MEANS JUST ONE REPLICA and the data could be lost if that replica is unavailable."

Topic with 3 replicas, min.insync.replicas = 2.

a replica falls behindanother falls behind3 replicas in synceverything proceeds normally2 replicas in synceverything proceeds normally1 replica in syncbrokers NO LONGER ACCEPT PRODUCE REQUESTS
  • · producers get NotEnoughReplicasException
  • · consumers can continue reading existing data
  • · ⇒ “a single in-sync replica becomes read-only”
Figure 7.3.83.3 min.insync.replicas — the fix for the silent-degradation hole

Why this is the right behavior:

"This prevents the undesirable situation where data is produced and consumed, only to disappear when unclean election occurs."

Recovery: "we must make one of the two unavailable partitions available again (maybe restart the broker) and wait for it to catch up and get in sync."

min.insync.replicas trades availability for honesty.

  • Without it: writes keep succeeding as durability silently collapses.
  • With it: writes fail loudly the moment durability collapses.

A failed produce you can retry is strictly better than an acknowledged write you will lose.

3.4 Keeping replicas in sync

Two configs, matching the two ways a replica goes out of sync:

ConfigGuards againstDefault historyTuning guidance
zookeeper.session.timeout.msLosing connectivity to ZooKeeper6 s → 18 s in 2.5.0, "in order to increase the stability of Kafka clusters in cloud environments where network latencies show higher variance""high enough to avoid random flapping caused by garbage collection or network conditions, but still low enough to make sure brokers that are actually frozen will be detected in a timely manner"
replica.lag.time.max.msFalling behind the leader10 s → 30 s in 2.5.0, "to improve resilience of the cluster and avoid unnecessary flapping"⚠️ "this higher value also impacts MAXIMUM LATENCY FOR THE CONSUMER — with the higher value it can take up to 30 seconds until a message arrives to all replicas and the consumers are allowed to consume it."

That last caveat is the hidden cost of the 2.5.0 default change. Raising replica.lag.time.max.ms buys cluster stability and pays with a worse worst-case consume latency — because visibility is gated on ISR replication (Ch. 6 §5.6).

3.5 Persisting to disk — and why Kafka mostly doesn't

"Kafka will acknowledge messages that were not persisted to disk, depending just on the number of replicas that received the message. Kafka will flush messages to disk when rotating segments (by default 1 GB in size) and before restarts but will otherwise rely on Linux page cache to flush messages when it becomes full."

The reasoning:

"having three machines in separate racks or availability zones, each with a copy of the data, is SAFER than writing the messages to disk on the leader, because simultaneous failures on two different racks or zones are so unlikely."

fsync on one leader3 copies in 3 racks/AZs (no fsync)
Survivespower loss on that hostpower loss on that host, disk failure, host failure, rack failure
Does not survivedisk failure, host losssimultaneous failure of 3 independent racks/AZs

Kafka bets on independence of failure domains rather than on the durability of a single device. That bet only pays off if your replicas really are in independent domains — hence broker.rack.

If you want fsync anyway:

  • flush.messages → maximum number of messages not synced to disk
  • flush.ms → frequency of syncing to disk

"Before using this feature, it is worth reading how fsync impacts Kafka's throughput and how to mitigate its drawbacks."


On this page