12.1 Transmitting Event Streams
An EVENT is the streaming counterpart of a batch record:
A small, SELF-CONTAINED, IMMUTABLE object containing the details of SOMETHING THAT HAPPENED AT A POINT IN TIME. It usually contains a TIMESTAMP indicating when it happened according to a TIME-OF-DAY clock.
Vocabulary mapping:
| Batch | Stream |
|---|---|
| file written once, read by multiple jobs | event generated once by a PRODUCER (publisher, sender), processed by multiple CONSUMERS (subscribers, recipients) |
| filename identifies a set of related records | TOPIC or STREAM groups related events |
Why not just poll a database?
In principle a file or database is sufficient: the producer writes every event, and each consumer PERIODICALLY POLLS for events since it last ran. THIS IS ESSENTIALLY WHAT A BATCH PROCESS DOES.
But when moving toward continual processing with low delays, POLLING BECOMES EXPENSIVE if the datastore isn't designed for it. THE MORE OFTEN YOU POLL, THE LOWER THE PERCENTAGE OF REQUESTS THAT RETURN NEW EVENTS, AND THUS THE HIGHER THE OVERHEADS.
Databases have traditionally not supported notification well. Relational databases have TRIGGERS, but they are VERY LIMITED and have been SOMEWHAT OF AN AFTERTHOUGHT in database design.
1.1 The two questions that differentiate every messaging system
1 · What if producers send faster than consumers can process? There are three options — drop messages, buffer them in a queue, or apply backpressure (flow control, blocking the producer).
- Unix pipes and TCP use backpressure: a small fixed-size buffer, and if it fills, the sender is blocked until the recipient takes data out.
- If buffering: “Does the system crash if the queue no longer fits in memory, or does it write messages to disk? In the latter case, how does the disk access affect performance, and what happens when the disk fills up?”
2 · What if nodes crash or go offline — are messages lost? “Durability may require a combination of writing to disk and/or replication, which has a cost. If you can afford to sometimes lose messages, you can probably get higher throughput and lower latency on the same hardware.”
Whether loss is acceptable is application-specific:
With periodic sensor readings and metrics, an occasional missing data point is perhaps not important — an updated value arrives shortly. HOWEVER, BEWARE THAT IF A LARGE NUMBER OF MESSAGES ARE DROPPED, IT MAY NOT BE IMMEDIATELY APPARENT THAT THE METRICS ARE INCORRECT.
If you are COUNTING events, reliable delivery matters more, since EVERY LOST MESSAGE MEANS INCORRECT COUNTERS.
1.2 Direct messaging — and why it's limited
| Approach | Where used |
|---|---|
| UDP multicast | Widely used in FINANCE for stock market feeds, where low latency is important. UDP is unreliable, but application-level protocols can recover lost packets — the producer must REMEMBER PACKETS IT HAS SENT so it can retransmit on demand |
| Brokerless libraries (ZeroMQ, nanomsg) | Publish/subscribe over TCP or IP multicast |
| StatsD (metrics agents) | Unreliable UDP. "Counter metrics are correct ONLY IF ALL MESSAGES ARE RECEIVED; using UDP makes the metrics AT BEST APPROXIMATE" |
| Webhooks | A callback URL of one service registered with another, which requests that URL whenever an event occurs |
These generally require THE APPLICATION CODE TO BE AWARE OF THE POSSIBILITY OF MESSAGE LOSS. The faults they tolerate are QUITE LIMITED — they generally assume PRODUCERS AND CONSUMERS ARE CONSTANTLY ONLINE.
If a consumer is OFFLINE, it may miss messages sent while unreachable. Some protocols let the producer retry — BUT THIS BREAKS DOWN IF THE PRODUCER CRASHES, LOSING THE BUFFER OF MESSAGES IT WAS SUPPOSED TO RETRY.
1.3 Message brokers
A kind of DATABASE OPTIMIZED FOR HANDLING MESSAGE STREAMS. By CENTRALIZING the data, these systems more easily tolerate clients that come and go, and THE QUESTION OF DURABILITY IS MOVED TO THE BROKER.
Faced with slow consumers, they generally allow UNBOUNDED QUEUEING (as opposed to dropping or backpressure).
A consequence of queueing is that CONSUMERS ARE GENERALLY ASYNCHRONOUS: when a producer sends, it normally waits ONLY for the broker to confirm it has BUFFERED the message — NOT for consumers to process it.
Brokers vs databases — four differences:
| Database | Message broker | |
|---|---|---|
| Retention | Keeps data until EXPLICITLY DELETED | Some AUTOMATICALLY DELETE a message once successfully delivered. NOT SUITABLE FOR LONG-TERM DATA STORAGE |
| Working set | Large | Assumes queues are SHORT. If it must buffer a lot because consumers are slow (spilling to disk), EACH MESSAGE TAKES LONGER TO PROCESS AND OVERALL THROUGHPUT MAY DEGRADE |
| Selection | Secondary indexes, a query language | Subscribing to a subset of topics matching a pattern — "both are ways for a client to select the portion of the data it wants, but databases offer MUCH MORE ADVANCED query functionality" |
| Change notification | Result is a POINT-IN-TIME SNAPSHOT; the client is NOT told when its result becomes outdated unless it repeats the query or polls | No arbitrary queries and no updates after send — BUT THEY NOTIFY CLIENTS WHEN DATA CHANGES |
(Standards: JMS, AMQP. Implementations: RabbitMQ, ActiveMQ, HornetQ, Qpid, TIBCO EMS, IBM MQ, Azure Service Bus, Google Cloud Pub/Sub. And: "although it is possible to use databases as queues, TUNING THEM TO GET GOOD PERFORMANCE IS NOT STRAIGHTFORWARD.")
Two consumer patterns — and their combination:
Kafka consumer groups combine the two: within a group you get load balancing, across groups you get fan-out.
1.4 Acknowledgments, redelivery, and the reordering trap
Brokers use ACKNOWLEDGMENTS: a client must explicitly tell the broker when it has FINISHED PROCESSING so the broker can remove the message. If the connection closes or times out without an ack, the broker ASSUMES THE MESSAGE WAS NOT PROCESSED and DELIVERS IT AGAIN to another consumer.
⚠️ "It could happen that the message ACTUALLY WAS FULLY PROCESSED, BUT THE ACKNOWLEDGMENT WAS LOST IN THE NETWORK. Handling this requires an ATOMIC COMMIT PROTOCOL — unless the operation was IDEMPOTENT or exactly-once semantics are not required."
Load balancing + redelivery ⇒ inevitable reordering:
The producer sends m1 m2 m3 m4 m5 in that order, to two consumers sharing the queue.
Consumer 1 processes m4, then m3 — not the order the producer sent them.
“Even if the broker otherwise tries to preserve order (as required by both JMS and AMQP), the combination of load balancing with redelivery inevitably leads to messages being reordered.”
- To avoid it, use a separate queue per consumer — that is, do not use load balancing.
- Not a problem if messages are independent; important if there are causal dependencies.
The poison-message loop and dead letter queues:
A common scenario: a producer IMPROPERLY SERIALIZES a message — e.g. leaving out a required key in JSON. The message causes a consumer to CRASH AND RESTART, so it won't acknowledge, so the broker RESENDS IT, which causes ANOTHER CONSUMER TO FAIL. THIS LOOP REPEATS ITSELF INDEFINITELY.
If the broker guarantees STRONG ORDERING, NO FURTHER PROGRESS CAN BE MADE. Brokers that allow reordering can continue — but WILL WASTE RESOURCES on messages that will never be acknowledged.
DEAD LETTER QUEUES (DLQs) fix this: rather than retrying forever, the message is MOVED TO A DIFFERENT QUEUE TO UNBLOCK CONSUMERS. MONITORING IS USUALLY SET UP ON DLQs — ANY MESSAGE IN THE QUEUE IS AN ERROR. An operator can then permanently drop it, manually modify and reproduce it, or fix consumer code to handle it.