Learn Labs
12. Stream Processing

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:

BatchStream
file written once, read by multiple jobsevent generated once by a PRODUCER (publisher, sender), processed by multiple CONSUMERS (subscribers, recipients)
filename identifies a set of related recordsTOPIC 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

ApproachWhere used
UDP multicastWidely 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"
WebhooksA 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:

DatabaseMessage broker
RetentionKeeps data until EXPLICITLY DELETEDSome AUTOMATICALLY DELETE a message once successfully delivered. NOT SUITABLE FOR LONG-TERM DATA STORAGE
Working setLargeAssumes 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
SelectionSecondary indexes, a query languageSubscribing 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 notificationResult is a POINT-IN-TIME SNAPSHOT; the client is NOT told when its result becomes outdated unless it repeats the query or pollsNo 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:

(a) Load balancing — each message to one consumer
one of themtopicconsumer 1consumer 2consumer 3Share the work; add consumers to parallelize when messages are expensive to process.AMQP: multiple clients on one queue. JMS: shared subscription.
(b) Fan-out — each message to all consumers
all of themtopicconsumer 1gets ALLconsumer 2gets ALLconsumer 3gets ALL“Several independent consumers each tune in to the same broadcast, without affecting one another”— the streaming equivalent of several batch jobs reading the same input file.JMS: topic subscriptions. AMQP: exchange bindings.

Kafka consumer groups combine the two: within a group you get load balancing, across groups you get fan-out.

Figure 12.1.2Two consumer patterns — and their combination

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.

brokerconsumer 1consumer 2m1m2m3 — crashes before acknowledgingm4unacknowledged m3, redelivered

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.
Figure 12.1.3Load balancing + redelivery ⇒ inevitable reordering

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.


On this page