Skip to content
LakeBench
ProblemsCommunityPricing
Sign inStart practicing

Streaming & Message Queues

Progress0/13
x

Batch versus Streaming

  • Batch versus streaming12m
  • Streaming architectures12m
  • Exactly-once semantics14m

Why Queues

  • Why a queue12m
  • Topics, partitions, keys14m

Offsets and Groups

  • Offsets and commits14m
  • Consumer groups14m

Delivery and Time

  • At-least-once14m
  • Late and out of order14m
  • Watermarks and Spark14m

Stream-batch

  • Stream-batch unification12m
  • Replay from offset12m

Capstone

  • Capstone: late messages16m
Back to track
  1. Learn
  2. Streaming & Message Queues
  3. Why Queues
  4. Why a queue

Lesson 4 of 13 · Theory first, then run it

Why a queue

streamingpythonbeginner12 min

Overview

A click should not wait for The night job. Decouple, buffer, replay: three jobs HTTP cannot own.

On this page5 sections›
  1. 1The idea
  2. 2Why this exists
  3. 3Common beginner questions
  4. 4What comes next
  5. 5Practice

The idea

A message queue is a buffer between a data producer and a data consumer. The producer writes messages (events, records, signals) into the queue and immediately moves on to its next task. The consumer reads messages from the queue at its own pace, whenever it is ready. The queue stores messages until they are consumed.

Think of it like a mailbox. The mail carrier (producer) drops off letters and leaves. You (consumer) check the mailbox whenever you want. The mailbox (queue) holds the letters in between. The carrier does not need to wait for you to be home, and you do not need to be home when the carrier arrives.

FakeKafkaTopic is a Python class that simulates a message queue using lists. No Kafka broker runs in this tab. You practice produce, consume, and replay on this object because the habits transfer to any real messaging system.

Why this exists

Imagine an e-commerce website that tracks every click. Each click generates an event that needs to be stored in the data warehouse. Without a queue, the click-tracking code must write directly to the warehouse. If the warehouse is slow or down, the user's page hangs waiting for the write to complete. With a queue, the click-tracking code drops the event into the queue (fast, always available) and the user's page loads instantly. A separate consumer reads events from the queue and writes them to the warehouse at its own pace.

A message queue provides three essential capabilities that direct HTTP calls cannot:

Producer, queue, consumer
Click event(producer)Message queue(buffer)WarehouseconsumerGold table

The producer drops a message and leaves. The consumer reads when ready. The queue holds everything in between.

  1. Decoupling: the producer does not wait for the consumer. If the warehouse is slow, the producer is unaffected.
  2. Buffering: if the consumer goes down, messages accumulate in the queue safely. Nothing is lost.
  3. Replay: unlike an HTTP request that disappears after delivery, messages stay in the queue log. You can rewind and re-read them.
CapabilityWithout a queueWith a queue
DecoupleClick page waits on warehouse writeClick page writes to queue and returns instantly
BufferWarehouse outage drops eventsMessages sit safely until the consumer is back
ReplayAsk the source to resend (often impossible)Seek to an offset and re-read from the log
Python
# FakeKafkaTopic: simulating a message queue as a Python list
class FakeKafkaTopic:
    def __init__(self):
        self.messages = []

    def produce(self, message):
        """Producer drops a message and returns immediately."""
        self.messages.append(message)

    def consume(self, offset, count):
        """Consumer reads 'count' messages starting from 'offset'."""
        return self.messages[offset:offset + count]

# Usage
topic = FakeKafkaTopic()
topic.produce({"event": "click", "user": "u1", "page": "/products"})
topic.produce({"event": "click", "user": "u2", "page": "/cart"})

# Consumer reads when ready (not immediately)
batch = topic.consume(offset=0, count=2)
print(batch)

def why_a_queue():
    return ["decouple", "buffer", "replay"]

result = why_a_queue()
print(result)

Common beginner questions

Why not just use HTTP?

An HTTP POST delivers a message immediately, but if the receiver is down, the message is lost. HTTP is synchronous: the sender waits for a response. A queue is asynchronous: the sender drops the message and moves on. Queues also support replay (re-reading old messages), which HTTP does not.

Is Kafka the only message queue?

No. RabbitMQ, AWS SQS, Google Pub/Sub, Azure Event Hubs, and Apache Pulsar are all message queue systems. Kafka is the most popular for data engineering because of its high throughput, durability, and replay capability. The concepts in this track apply to all of them.

Does the queue ever run out of space?

In production, queues have retention policies. Kafka keeps messages for a configurable period (often 7 days). After that, old messages are deleted. The consumer must read within the retention window or lose access to old data.

The cluster is outside this platform

FakeKafkaTopic is the object model you will recognize on a real Kafka topic. Passing this exercise does not start Kafka. Installing and operating a broker is a step you take outside this platform.

A queue is not a database

Queues are for transporting messages, not for permanent storage. Consumers must process messages and write results to a durable store (like the warehouse). Do not use a queue as your data warehouse.

What comes next

A single queue (topic) works for low traffic. For high throughput, you need to split the topic into partitions so multiple consumers can read in parallel. That is the next lesson.

Practice

Return ["decouple", "buffer", "replay"] from why_a_queue and store that list in result.

Practicals · load into the editor

After you read the theory, run these in the pane on the right. They execute in this tab, no cluster.

Rate:
Was this useful?
Exactly-once semanticsTopics, partitions, keys