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

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."
  • 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 MillWhell diagram Example from MillWheel