1.3 How does it work internally? (the core object model)
A message is the unit of data — think a DB row or record.
3.1 Message and batch
A message is the unit of data — think a DB row or record. To Kafka it is an opaque array of bytes. Kafka does not know or care what's inside.
A message optionally has a key — also an opaque byte array. The key exists for partition routing control:
partition = hash(key) mod num_partitionsConsequence, and it's a big one: messages with the same key always land in the same partition — provided the partition count does not change. That parenthetical is a production landmine covered in §6.
Batching. Messages are written in batches — a batch is a collection of messages all bound for the same topic and the same partition. Why: a network round trip per message is unacceptable overhead.
| Without batching | With batching |
|---|---|
msg → RTTmsg → RTTmsg → RTT… | [msg,msg,msg,...,msg] → 1 RTT |
| A network round trip per message — unacceptable overhead. | Every message in the batch is bound for the same topic AND the same partition, and the batch is typically compressed. |
The tradeoff is explicit: latency vs throughput.
- Larger batches → more messages/unit time, but longer for any individual message to propagate.
- Batches are typically compressed → cheaper transfer and storage, at the cost of CPU.
3.2 Schemas — Kafka's deliberate omission
Kafka treats payloads as bytes. The book's position is that you should nonetheless impose a schema, and it explains why in terms of decoupling deploys, not tidiness.
| Option | Pros | Cons |
|---|---|---|
| JSON / XML | Easy, human-readable | No robust type handling; no cross-version compatibility |
| Avro (favored) | Compact; schema separate from payload; no codegen needed when schema changes; strong typing; backward and forward compatible evolution | Needs a schema repository |
Why this is an architecture concern, not a serialization preference:
Without schemas, writing and reading are tightly coupled, which forces this deployment dance:
- All consumers must be updated to handle Both old and new formats
- … and deployed
- Only then can producers be updated to emit the new format
A lockstep, ordered, cross-team deploy for every field addition — exactly the failure that made LinkedIn's XML tracking system break "constantly." With well-defined schemas in a common repository, messages are understandable without coordination. Schema registry is how you buy back independent deployability.
3.3 Topics and partitions
Topic ≈ a database table or a filesystem folder. Partition = a single log.
Ordering is guaranteed within a partition only. There is NO ordering guarantee across a topic.
The single most-violated fact in Kafka:
Ordering is guaranteed within a partition only. There is NO ordering guarantee across a topic.
If your correctness depends on global ordering, you need either one partition (throughput ceiling of a single broker/consumer) or a key scheme that co-locates all causally-related events in one partition.
Partitions are the mechanism for both scale and redundancy:
- Scale: each partition can live on a different server, so one topic scales horizontally far beyond a single machine.
- Redundancy: partitions can be replicated across servers, so a copy survives a server failure.
"Stream" — usually means a single topic of data, regardless of partition count — a single flow producers→consumers. The term is most used in stream processing (Kafka Streams, Samza, Storm) which acts on messages in real time, versus offline frameworks (Hadoop) that act on bulk data later.
3.4 Producers and consumers
Client tiers:
The advanced clients are built on producers and consumers. There is no separate protocol underneath — a useful thing to remember when debugging Connect or Streams.
The advanced clients are built on producers and consumers. There is no separate protocol underneath — a useful thing to remember when debugging Connect or Streams.
Producers (elsewhere: publishers/writers) create messages to a specific topic.
- Default: balance messages evenly over all partitions.
- With a key: hash the key → map to a partition → all messages with that key go to one partition.
- Custom partitioner: arbitrary business rules for message→partition mapping.
Consumers (elsewhere: subscribers/readers) subscribe to topics and read messages in produced order per partition.
Offsets. Kafka adds an offset to each message at produce time — an integer that continually increases. Each message in a partition has a unique offset; the next message has a greater offset — though not necessarily monotonically greater (i.e. gaps are legal; do not write code that assumes offset+1). By storing the next offset per partition — typically inside Kafka itself — a consumer can stop and restart without losing its place.
Consumer groups.
- INVARIANT: each partition is consumed by exactly ONE member of the group. A single consumer may own several partitions — that's “ownership”.
- Two properties fall out: horizontal consumer scaling (add members to divide the partitions) and automatic failover (if a member dies, the remaining members reassign its partitions). The ceiling: group parallelism ≤ partition count.
Two properties fall out:
- Horizontal consumer scaling — add members to divide the partitions.
- Automatic failover — if a member dies, the remaining members reassign its partitions.
The ceiling this implies: group parallelism ≤ partition count. Extra consumers beyond the partition count sit idle. This is why partition count is a capacity decision made early (see Ch. 2).
3.5 Brokers and clusters
A broker is a single Kafka server. Its job:
- Receive messages from producers
- Assign offsets to them
- Write messages to storage on disk
- Service consumer fetch requests
Capacity of one broker (hardware-dependent): thousands of partitions and millions of messages/second.
Cluster + controller.
- The controller is elected AUTOMATICALLY from the live cluster members.
- Responsibilities: assign partitions to brokers · monitor for broker failures.
Leaders and followers.
- A partition is owned by a single broker — that broker is the leader. A replicated partition is also assigned to other brokers, its followers.
- All producers must connect to the leader to publish. Consumers may fetch from the leader or a follower. If the leader's broker fails, a follower takes over leadership.
- A partition is owned by a single broker → that broker is the leader.
- A replicated partition is also assigned to other brokers → followers.
- All producers must connect to the leader to publish.
- Consumers may fetch from the leader OR a follower.
- If the leader's broker fails, a follower takes over leadership.
3.6 Retention — the feature that changes everything
Retention = durable storage of messages for a period of time. Brokers have a default; topics can override.
Two policies:
- By time — e.g. 7 days
- By size — e.g. 1 GB per partition
When a limit is hit, messages are expired and deleted. So retention configuration defines a minimum amount of data available at any time.
Per-topic tuning is the point: a tracking topic might keep several days; application metrics only a few hours.
Log compaction — a third mode: retain only the last message produced with a specific key. Purpose: changelog-type data, where only the latest update matters. This is what makes Kafka usable as a state store rather than only a pipe (see Ch. 14 on KTables).
Both topics receive the same four records, in this order:
k=A,v=1
k=B,v=1
k=A,v=2
k=A,v=3| Normal retention (delete) | Log compaction |
|---|---|
| All four are kept until the time or size limit is hit, then expired and deleted. | The earlier values for a key are superseded, eventually removed. Final: k=B,v=1 and k=A,v=3 — the last value per key survives. |
| Retention by time (e.g. 7 days) or by size (e.g. 1 GB per partition). | For changelog-type data, where only the latest update matters — what makes Kafka usable as a state store rather than only a pipe. |
3.7 Multiple clusters and MirrorMaker
Three reasons deployments grow to multiple clusters:
- Segregation of types of data
- Isolation for security requirements
- Multiple datacenters (disaster recovery)
Critical constraint: "The replication mechanisms within the Kafka clusters are designed only to work within a single cluster, not between multiple clusters."
Cross-cluster copying uses MirrorMaker, which is architecturally almost insultingly simple: a Kafka consumer and a Kafka producer, linked together with a queue. Consume from cluster A, produce to cluster B.
- MirrorMaker is a Kafka consumer and a Kafka producer, linked together with a queue — consume from cluster A, produce to cluster B.
- It exists because “the replication mechanisms within the Kafka clusters are designed only to work within a single cluster, not between multiple clusters.”
Motivating example from the book: a user edits public profile info; that change must be visible regardless of which datacenter serves the search results. Or: collect monitoring data from many sites into one central place where analysis/alerting lives.
"The simple nature of the application belies its power in creating sophisticated data pipelines." Full treatment in Ch. 9/10.