12.2 Log-Based Message Brokers
The mindset difference this fixes:
| AMQP/JMS-style | Databases and filesystems |
|---|---|
| Inherited the transient messaging mindset of network packets: even if written to disk, messages are quickly deleted after delivery. | “Everything written is normally expected to be permanently recorded, at least until someone explicitly chooses to delete it again.” |
- “Receiving a message is destructive if the acknowledgment causes it to be deleted, so you cannot run the same consumer again and expect the same result.”
- “If you add a new consumer, it typically receives only messages sent after it was registered; any prior messages are already gone and cannot be recovered.”
The hybrid is to combine the durable storage approach of databases with the low-latency notification facilities of messaging.
2.1 How it works
A log is simply an append-only sequence of records on disk — we have already seen it as LSM storage and the WAL in Chapter 4, as replication in Chapter 6, and as consensus in Chapter 10. The producer appends to the end of the log; the consumer reads the log sequentially and, at the end, waits for a notification. “The Unix tool tail -f essentially works like this.”
To scale beyond one disk, shard the log (Chapter 7). Each shard is a separate log, read and written independently, and a topic is a group of shards carrying messages of the same type.
Kafka, Amazon Kinesis Streams work this way; Google Cloud Pub/Sub is architecturally similar but exposes a JMS-style API rather than a log abstraction.
Even though these brokers write ALL messages to disk, they achieve throughput of MILLIONS OF MESSAGES PER SECOND by sharding across machines, and fault tolerance by replicating.
2.2 Log vs traditional messaging — the honest comparison
Fan-out is trivial: "several consumers can independently read the log without affecting one another; READING A MESSAGE DOES NOT DELETE IT."
Load balancing is coarse: "the broker assigns ENTIRE SHARDS to nodes in the consumer group instead of assigning individual messages."
Two downsides of coarse-grained balancing:
- "The number of nodes sharing the work can be AT MOST THE NUMBER OF LOG SHARDS in that topic." (You could have two consumers split even/odd offsets, or use a thread pool — but that complicates consumer offset management. In general, SINGLE-THREADED PROCESSING OF A SHARD IS PREFERABLE, AND PARALLELISM CAN BE INCREASED BY USING MORE SHARDS.)
- "If a SINGLE MESSAGE IS SLOW TO PROCESS, IT HOLDS UP THE PROCESSING OF SUBSEQUENT MESSAGES IN THAT SHARD" — head-of-line blocking (Ch 2).
THE DECISION RULE: when messages may be EXPENSIVE to process, you want to parallelize PER MESSAGE, and ORDERING IS NOT SO IMPORTANT → JMS/AMQP style. When throughput is HIGH, each message is FAST to process, and ORDERING IS IMPORTANT → LOG-BASED.
(The distinction is blurring — Kafka now supports JMS/AMQP-style consumer groups allowing multiple consumers to receive messages from the same partition.)
Routing for ordering: "Since sharded logs preserve ordering only WITHIN a shard, ALL MESSAGES THAT NEED TO BE CONSISTENTLY ORDERED MUST BE ROUTED TO THE SAME SHARD. E.g. events relating to one particular user appear in a fixed order — achieved by making THE USER ID THE PARTITION KEY."
2.3 Consumer offsets
Consuming a shard sequentially makes it easy to tell which messages have been processed: all messages with an offset LESS than the current offset are done. THE BROKER DOES NOT NEED TO TRACK ACKNOWLEDGMENTS FOR EVERY MESSAGE — ONLY TO PERIODICALLY RECORD THE CONSUMER OFFSETS. The reduced bookkeeping and the opportunities for BATCHING AND PIPELINING help increase throughput.
⚠️ "If a consumer fails, it will RESUME FROM THE LAST RECORDED OFFSET rather than the more recent last offset it saw. THIS CAN CAUSE THE CONSUMER TO SEE SOME MESSAGES TWICE."
The offset is in fact very similar to the LOG SEQUENCE NUMBER in single-leader replication. EXACTLY THE SAME PRINCIPLE: THE MESSAGE BROKER BEHAVES LIKE A LEADER DATABASE AND THE CONSUMER LIKE A FOLLOWER.
2.4 Disk space, and the ring buffer
The log is divided into Segments; old segments are Deleted or Moved to archive.
⇒ “Effectively, the log implements a Bounded-size buffer that discards old messages when it gets full — a Circular buffer or Ring buffer. However, since that buffer is On disk, it can be quite large.”
Back-of-the-envelope:
- 20 TB drive, 250 MB/s sequential write
20e12 / 250e6 = 80,000 s ≈ 22 HOURSto fill at the Fastest possible rate.
⇒ “A disk-based log can Always buffer At least 22 hours’ worth of messages, even with many disks and machines (more disks increases Both space And total write bandwidth). In practice, deployments rarely use the full write bandwidth, so the log can typically keep Several days’ or even weeks’ worth.”
Tiered and object storage: Kafka and Redpanda serve older messages from object storage as TIERED STORAGE. WarpStream, Confluent Freight, and Bufstream store ALL data in the object store.
In addition to COST EFFICIENCY, this architecture MAKES DATA INTEGRATION EASIER: messages in object storage are stored AS ICEBERG TABLES, which enable BATCH AND DATA WAREHOUSE JOB EXECUTION DIRECTLY ON THE DATA WITHOUT HAVING TO COPY IT INTO ANOTHER SYSTEM.
When consumers can't keep up — the operational advantages:
The log-based approach is Buffering with a Large but fixed-size buffer. If a consumer falls behind the retention window, It misses messages — the broker effectively Drops old messages.
- ✔ “You can Monitor how far a consumer is behind the head of the log and raise an alert. As the buffer is large, there is enough time for a human operator to fix the slow consumer before it starts missing messages.”
- ✔ “Even if a consumer Does fall too far behind, Only that consumer is affected; it does not disrupt the service for other consumers. this is a big operational advantage: you can Experimentally consume a production log for development, testing, or debugging Without having to worry much about disrupting production.”
- ✔ “When a consumer is shut down or crashes, It stops consuming resources — the only thing that remains is its consumer offset.”
- ✗ Contrast with traditional brokers: “you need to be careful to Delete any queues whose consumers have been shut down, to avoid them unnecessarily Accumulating messages and taking away memory from active consumers.”
2.5 Replaying old messages — the batch-like property
In a log-based broker, consuming messages is MORE LIKE READING FROM A FILE: a READ-ONLY operation that does not change the log. The only side effect is that THE CONSUMER OFFSET MOVES FORWARD — and THE OFFSET IS UNDER THE CONSUMER'S CONTROL.
You can start a copy of a consumer WITH YESTERDAY'S OFFSET and write output to a DIFFERENT LOCATION to reprocess the last day's messages. YOU CAN REPEAT THIS ANY NUMBER OF TIMES, VARYING THE PROCESSING CODE.
This makes log-based messaging MORE LIKE THE BATCH PROCESSES OF CH 11, where derived data is clearly separated from input data through a REPEATABLE TRANSFORMATION PROCESS. It allows more experimentation and easier recovery from errors and bugs — A GOOD TOOL FOR INTEGRATING DATAFLOWS WITHIN AN ORGANIZATION.