Integer identifier, default 0, must be unique per broker within a cluster.
The selection is technically arbitrary and can be moved between brokers for maintenance. But:
Set it to something intrinsic to the host so that during maintenance it is not onerous to map broker IDs to hosts. If hostnames are host1.example.com, host2.example.com, use 1 and 2.
Why this matters at 3am: every metric, log line, and --describe output refers to brokers by ID. If ID↔host requires a lookup table, you're slower during an incident.
hostname / port — ZooKeeper server and its client port.
/path — optional chroot to use as the root for this Kafka cluster. Omitted → root path. If the chroot path doesn't exist, the broker creates it at startup.
It lets the ZooKeeper ensemble be shared with other applications, including other Kafka clusters, without conflict.
Also: specify multiple ZooKeeper servers (all in the same ensemble) so the broker can connect to another member on server failure.
Practical consequence: if you ever want a second Kafka cluster (and Ch. 1 says you eventually will — for data segregation, security isolation, or DR), a chroot from day one means you don't have to migrate ZooKeeper state later.
Where log segments live. log.dirs is a comma-separated list of local paths; if unset it falls back to log.dir.
Placement algorithm — and its flaw:
The broker stores partitions across directories in a "least-used" fashion, with one partition's log segments stored within the same path.
The broker places a new partition in the path that has the least number of partitions currently stored in it — NOT the least amount of disk space used. So an even distribution of data across multiple directories is not guaranteed.
log.dirs=/data1,/data2,/data3
Directory
What is in it
Partitions
/data1
P1 — 800 GB · P4 — 2 GB
2
/data2
P2 — 1 GB · P5 — 1 GB · P6 — 1 GB
3
/data3
P3 — 1 GB · P7 — 1 GB
2
The broker places a new partition in the path that has the least number of partitions currently stored in it — not the least amount of disk space used. So the next partition goes to /data1 or /data3, even though /data1 is nearly full. This is how you fill a disk. “An even distribution of data across multiple directories is not guaranteed.”
A configurable thread pool used in exactly three situations:
Normal startup — open each partition's log segments
Startup after a failure — check and truncate each partition's log segments
Shutdown — cleanly close log segments
Default: one thread per log directory.
Since these threads are only used during startup and shutdown, it is reasonable to set a larger number to parallelize. Specifically, when recovering from an unclean shutdown, this can mean the difference of several hours when restarting a broker with a large number of partitions!
Multiplication trap: the value is per log directory. num.recovery.threads.per.data.dir=8 with 3 paths in log.dirs = 24 threads total.
This can be undesirable, especially as there is no way to validate the existence of a topic through the Kafka protocol without causing it to be created.
That's a genuinely nasty property: checking creates. Set to false if you manage topic creation explicitly (manually or via provisioning).
These are cluster-wide defaults for newly created topics. Set them to "baseline values appropriate for the majority of the topics in the cluster." Per-topic overrides go through the admin tools (Ch. 12).
Removed in newer versions:log.retention.hours.per.topic, log.retention.bytes.per.topic, log.segment.bytes.per.topic — these old broker-side per-topic overrides are no longer supported. Use admin tools.
Partitions for a new topic, primarily when auto-creation is enabled. Default: 1.
Partition count can only be increased, never decreased. If a topic needs fewer partitions than num.partitions, you must create it manually.
Many users set partition count equal to, or a multiple of, the number of brokers so partitions distribute evenly → message load distributes evenly. A 10-partition topic on a 10-broker cluster with balanced leadership has optimal throughput. Not a requirement — you can balance load other ways, e.g. multiple topics.
What throughput do you expect for the topic? (100 KBps? 1 GBps?)
What is your maximum throughput consuming from a single partition?A partition is always consumed completely by a single consumer — even without consumer groups, the consumer must read all messages in the partition. If your slow consumer writes to a database that never handles more than 50 MBps per writing thread, you are limited to 50 MBps per partition.
You can do the same exercise for per-producer per-partition throughput, but producers are typically much faster than consumers, so it's usually safe to skip.
If you send messages to partitions based on keys, adding partitions later can be very challenging → calculate throughput on expected future usage, not current.
Consider partitions per broker, and available disk space and network bandwidth per broker.
Avoid overestimating — each partition uses memory and other resources on the broker and increases the time for metadata updates and leadership transfers.
Are you mirroring data? Factor in mirroring throughput. Large partitions can become a bottleneck in many mirroring configurations.
Cloud IOPS limits. There may be hard IOPS caps per VM/disk. Too many partitions increases IOPS due to the parallelism involved → you hit quotas.
The denominator is the consumer, not the producer: a partition is always consumed completely by a single consumer, so if your slow consumer writes to a database that never handles more than 50 MBps per writing thread, you are limited to 50 MBps per partition. Producers are typically much faster than consumers, so that side is usually safe to skip.
Heuristic when you have no numbers:
Limit the size of the partition on disk to less than 6 GB per day of retention. Starting small and expanding as needed is easier than starting too large.
Summary of the tension:"you want many partitions, but not too many."
The recommendation, and the reasoning — this is one of the best passages in the chapter:
Set replication factor to at least 1 above min.insync.replicas.
For more fault resistance, if you have large enough clusters and enough hardware, set it 2 above min.insync.replicas — abbreviated RF++.
RF++ allows easier maintenance and prevents outages.
Why RF++:"to allow for one planned outage within the replica set and one unplanned outage to occur simultaneously."
Set the replication factor at least 1 above min.insync.replicas; with large enough clusters and enough hardware, set it 2 above — RF++ — “to allow for one planned outage within the replica set and one unplanned outage to occur simultaneously.” For a typical cluster that means a minimum of three replicas of every partition.
Figure 2.4.5·default.replication.factor
For a typical cluster this means a minimum of three replicas of every partition. The scenario to hold in your head: a network switch outage, a disk failure, or some other unplanned problem during a rolling deployment or upgrade of Kafka or the OS. Rolling upgrades are exactly when your redundancy budget is already partly spent.
Time-based retention examines the last modified time (mtime) of each log segment file on disk. Under normal operations that's when the segment was closed — i.e. the timestamp of the last message in the file.
However, when using administrative tools to move partitions between brokers, this time is NOT accurate and will result in excess retention for these partitions. (Ch. 12 covers partition moves.)
Why: a partition move copies files, which resets mtime to "now" — so a 6-day-old segment looks brand new and survives another full retention period. Disk usage jumps after a rebalance and nobody knows why.
A topic with 8 partitions and log.retention.bytes = 1 GB retains at most 8 GB for the topic. All retention is performed for individual partitions, not the topic — so if you expand the partition count of a topic, retention also increases. -1 means infinite retention.
All retention is performed for individual partitions, not the topic. So if you expand the partition count of a topic, retention also increases when using log.retention.bytes.
-1 = infinite retention.
If both log.retention.bytes and a time parameter are set, messages may be removed when either criterion is met.
retention.ms = 1 day, retention.bytes = 1 GB: messages less than 1 day old can be deleted if the day's volume exceeds 1 GB.
Conversely, if volume is under 1 GB, messages are deleted after 1 day anyway.
Recommendation: for simplicity choose either size- or time-based retention — not both — to prevent surprises and unwanted data loss. Both can be used for advanced configurations.
Retention operates on log segments, not individual messages. Messages append to the current (active) segment. When it reaches log.segment.bytes (default 1 GB), it is closed and a new one opened. Only a closed segment can be considered for expiration.
Smaller segments → files closed and allocated more often → reduces overall efficiency of disk writes.
A topic receiving only 100 MB/day with default log.segment.bytes (1 GB) takes 10 days to fill one segment. Messages cannot expire until the segment is closed. If log.retention.ms = 1 week, there will actually be up to 17 days of messages retained: once the segment closes with 10 days of messages inside, that segment must be retained 7 more days before it expires — because the segment can't be removed until its last message can be expired.
A topic receiving only 100 MB/day with the default 1 GB log.segment.bytes takes 10 days to fill one segment, and only a closed segment can be considered for expiration — so the segment must then be retained 7 more days, “because the segment can’t be removed until its last message can be expired.” Fix: lower log.segment.bytesor set log.roll.ms for low-volume topics.
Figure 2.4.7·The low-volume-topic retention bug — work through this one
Fix: for low-volume topics, lower log.segment.bytesor set log.roll.ms. Otherwise your "1 week retention" is a lie and your GDPR/compliance deletion window is wrong.
When you request the offset for a partition at a specific timestamp, Kafka finds the log segment that was being written at that time — using the file's creation and last-modified times, looking for a file created before the timestamp and last modified after it. The offset at the beginning of that segment (which is also the filename) is returned.
Consequence: coarser segments = coarser timestamp seeks. With 1 GB segments on a slow topic, "seek to 10:15am" can land you days earlier. This directly affects --from-timestamp resets and time-based reprocessing.
Time after which a segment should be closed. No default — so by default segments close by size only.
log.segment.bytes and log.roll.ms are not mutually exclusive: Kafka closes a segment when either the size limit or the time limit is reached, whichever comes first.
When using a time-based segment limit, consider the impact of many log segments being closed simultaneously. This happens when many partitions never reach the size limit: the time-limit clock starts when the broker starts, so it always fires at the same moment for all those low-volume partitions.
The time-limit clock starts when the broker starts, so it always fires at the same moment for all those low-volume partitions. log.segment.bytes and log.roll.ms are not mutually exclusive — Kafka closes a segment when either limit is reached, whichever comes first.
Setting min.insync.replicas = 2 ensures at least two replicas are caught up and "in sync" with the producer. Used in tandem with the producer config acks=all. This ensures at least two replicas (leader + one other) acknowledge a write for it to be successful.
The failure it prevents:
min.insync.replicas = 2 ensures “at least two replicas are caught up and ‘in sync’ with the producer,” used in tandem with the producer config acks=all. The trade-off, stated plainly: higher durability is less efficient due to the extra overhead — clusters with high throughput that can tolerate occasional message loss aren’t recommended to change this from the default of 1.
Figure 2.4.9·The failure it prevents
Trade-off, stated plainly: configuring for higher durability is less efficient due to the extra overhead. Clusters with high throughput that can tolerate occasional message loss aren't recommended to change this from the default of 1.
That's an honest statement rarely made in docs: min.insync.replicas=1 is a deliberate data-loss-tolerant setting, and it's the default. (Full treatment in Ch. 7.)
Max size of a producible message. Default 1000000 (≈1 MB). A producer exceeding it gets an error back and the message is not accepted.
As with all byte sizes on the broker, this deals with compressed message size — producers can send messages much larger uncompressed, provided they compress to under the limit.
Performance impacts of raising it:
Broker threads handling network connections/requests work longer on each request
Larger disk writes → impacts I/O throughput
Alternatives mentioned: blob stores and/or tiered storage (not covered in the chapter)
broker : message.max.bytes consumer : fetch.message.max.bytes ← must be ≥ broker's brokers : replica.fetch.max.bytes ← must be ≥ broker's (for replication)
If fetch.message.max.bytes < message.max.bytes, consumers hitting a larger message fail to fetch it, resulting in a consumer that gets stuck and cannot proceed. Same rule for replica.fetch.max.bytes on brokers — otherwise replication stalls on that partition.
This is a "raise one number, break the cluster silently" trap. Raise all three together, in the right order (consumers/replicas first, broker last).