ddia/ch11-stream.md - DDIA
General
- Events are organized in topics or streams
- Polling is bad at irregular message arrival + granular level, because many poll requests return nothing.
- Disk: you can offload old data to object storage
Protocols/styles
- AMQP (Advanced Message Queuing Protocol)
- open standard for sending messages in middleware
- defining features of AMQP are message orientation, queuing, routing (including point-to-point and publish-and-subscribe), reliability and security.
- JMS (Java Message Service)
- previous tech before AMQP
- Both of the above delete messages after processing. That's kind of mental model there (as opposed to log-based message brokers)
- Log-based brokers
- Each message is appended to a log.
- Log can be sharded
- Each message within a partition has a monotonically increasing sequence number (called offset).
- Broker remembers consumer offsets for fault tolerance
- If a consumer node fails, another node in the consumer group is assigned the failed consumer’s shards, and it starts consuming messages at the last recorded offset
- What to choose?
- When messages may be expensive to process and you want to parallelize processing on a message-by-message basis, and message ordering is not so important, the JMS/AMQP style of message broker is preferable.
- Because JMS/AMQP-style brokers let any free worker grab the next individual message off the queue (true per-message load balancing), whereas log-based systems restrict parallelism to whole shards
- When each message is fast to process and message ordering is imporant, log-based is better
- the distinction is a bit blurred - Kafka supports AMQP/JMS-style consumer groups
- When messages may be expensive to process and you want to parallelize processing on a message-by-message basis, and message ordering is not so important, the JMS/AMQP style of message broker is preferable.
What to do when consumers don't keep up?
- Drop messages.
- Use backpressure (flow control)
- Sender is blocked until consumer is OK
- Example: TCP or Unix pipe
- Save messages in a queue.
Direct messaging
- UDP multicast
- Low-latency
- Widely used in finance for delivering market feeds
- ZeroMQ
- Implement pub/sub using TCP or IP multicast
- HTTP/RPC requests
- use webhooks to directly push messages to consumer
Message brokers
- Also called message queue
- Server with clients (producers and consumers)
- Async
- Two styles: AMQP/JMS vs. Log-based
How to use multiple consumers in message broker?
- Load balancing
- each message is delivered to only 1 consumer
- here, you must accept that messages can be delivered out-of-order
- Fan-out
- each message is delivered to ALL consumers
- In Kafka, consumer group guarantees delivery to one of the consumers there. If multiple consumer groups are subscribed to a topic, the message is delivered to all of the consumer groups (fan-out) across consumer groups.
Stream analytics
- Trick: use probabilistic data structures over the stream to get some insight on-demand.
- IVM (Incremental View Maintanence)
Dealing with functions of time in streams
- Distinguish between:
- Processing time = local system clock on the processing machine
- Event time = when the event happened
- Windowing by processing time introduces artifacts due to variations in processing rate. ![[Pasted image 20260711161311.png]]
- Types of windows
- Tumbling window = fixed length, non-overlapping with other windows
- Hopping window = fixed length, overlaps with other windows partially. Example: a five-minute window with a hop size of one minute would contain the events from 10:03:00 to 10:07:59, then the next window would cover events from 10:04:00 to 10:08:59.
- Sliding window = queue of events in time interval. "last 5 min". fixed length.
- Session window = NO fixed length. windows are separated by inactivity.
Watermarks
- If you do window over a stream of data, how do you know all events arrived? Some events could be just delayed but their event time is in your window
- Watermarking is a term for knowing when you stop waiting for new events
- watermark is a timestamp,
T, that flows through the stream alongside the actual data, carrying the assertion: "I believe no more events with an event-time earlier than T will arrive from now on."
- watermark is a timestamp,
- Types
- Bounded out-of-orderness
- Punctuated: the source itself embeds markers in the stream indicating "everything before this point is final" (e.g., a heartbeat from an upstream system that knows its own completeness).
Anomaly detection using streams
Example from MillWheel