Designing Data-Intensive Applications
Ch. 11

Transmitting Event Streams

Event streams deliver records continuously — via message brokers, partitioned logs, or direct sockets.

Unlike request-response, event streams deliver a sequence of records over time. Message brokers decouple producers and consumers. Partitioned logs add durability and replay — consumers can re-read history at their own pace.

Diagram
In practice

Apache Kafka is the default partitioned log for microservices — order events, audit logs, and CDC streams. RabbitMQ and AWS SQS route task queues with at-least-once delivery. For browser real-time updates, WebSockets (bidirectional) and Server-Sent Events (server push over HTTP) complement the async backend. Kafka transactions plus idempotent producers achieve exactly-once processing semantics within a cluster.

WhatsApp at scale

Message delivery events append to a partitioned log so multiple consumer groups — push notifications, analytics, moderation — read the same stream at independent offsets without competing for messages.

typescript — Kafka consumer group
// LinkedIn / Uber-style Kafka consumer with offset tracking
async function consumeMessages(groupId: string) {
  const consumer = kafka.consumer({ groupId });
  await consumer.subscribe({ topic: "user-events" });

  await consumer.run({
    eachMessage: async ({ partition, message }) => {
      const event = JSON.parse(message.value!.toString());
      await processEvent(event); // must be idempotent for at-least-once
      // offset committed after processing (or in transaction with side effect)
    },
  });
}
Diagram
Key Takeaways
  • Message brokers (RabbitMQ) route messages to consumers with optional persistence.
  • Partitioned logs (Kafka) retain messages for replay and multiple consumer groups.
  • Producers append; consumers track their offset in the log.
  • Backpressure and buffering handle speed mismatches between producers and consumers.
  • Exactly-once processing semantics require idempotent consumers, transactional commits, or deduplication.
  • Kafka, RabbitMQ, AWS SQS/SNS, WebSockets, and SSE cover async and real-time delivery.
KafkaRabbitMQWebSocketSSEpartitioned logconsumer groupoffset