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. Topics, partitions, keys

Lesson 5 of 13 · Theory first, then run it

Topics, partitions, keys

streamingpythonbeginner14 min

Overview

A topic is a named log. Partitions are aisles. The same key always hashes to the same aisle.

On this page7 sections›
  1. 1The picture
  2. 2Why this shape
  3. 3Walk the boxes
  4. 4A matching example
  5. 5Common beginner questions
  6. 6What comes next
  7. 7Practice

The picture

Same key always goes to the same partition
P0Message key: ORD-1P1Hash: sum of character codes = 323P2323 % 3 partitions = 2P3Append to partition 2

ORD-1 hashes to partition 2 every time. All events for ORD-1 are ordered within partition 2.

Distributed Kafka event pipeline showing producers, broker topic partitions, and consumer groups
Distributed event streaming pipeline: producers publish events to partitioned topics on Kafka brokers, and consumer groups scale out parallel reads with offset tracking.
Source: Wikimedia Commons (MediaWiki Architecture)CC BY-SA 4.0

A partition is a subdivision of a topic. Instead of one giant list of messages, a topic is split into multiple smaller lists (partitions). Each partition is an independent, ordered sequence of messages. Kafka uses partitions to achieve parallelism: multiple consumers can read different partitions simultaneously.

When a producer sends a message, it includes a key (like order_id or user_id). The key determines which partition the message goes to. All messages with the same key always go to the same partition, which guarantees ordering for that key.

FakeKafkaTopic in this lesson is a Python list of lists: one inner list per partition. No Kafka broker runs here, but the partition assignment logic is the same.

Why this shape

Imagine a topic receiving 100,000 click events per second. A single consumer cannot process all 100,000 per second. With 10 partitions, you can have 10 consumers each handling 10,000 events per second. Partitions are how streaming systems scale horizontally.

The key-to-partition assignment also matters for ordering. If all events for order ORD-1 go to partition 2, they are guaranteed to be in the order they were produced. If they were spread across partitions, you would lose ordering guarantees. A payment event might appear before the order creation event, causing errors.

Walk the boxes

Partition assignment works by hashing the message key and taking the modulo of the number of partitions. This ensures the same key always lands in the same partition.

A matching example

Run the example below in this tab. Read the input, follow the code, then check the output matches what you expect.

Python
# Partition assignment: deterministic hash of the key
def partition_for(key, n_partitions):
    """Assign a key to a partition using a stable hash."""
    # sum of ordinals: deterministic across runs (unlike hash())
    hash_value = sum(ord(c) for c in key)
    return hash_value % n_partitions

# ORD-1: O=79, R=82, D=68, -=45, 1=49 => 323
# 323 % 3 = 2
print(partition_for("ORD-1", 3))  # 2
print(partition_for("ORD-2", 3))  # (79+82+68+45+50) % 3 = 324 % 3 = 0
print(partition_for("ORD-3", 3))  # (79+82+68+45+51) % 3 = 325 % 3 = 1

Each key consistently maps to the same partition, spreading load across all partitions:

KeyOrdinal sumPartition (mod 3)
ORD-13232
ORD-23240
ORD-33251
Python
# FakeKafkaTopic with partitions
class FakeKafkaTopic:
    def __init__(self, n_partitions):
        self.partitions = [[] for _ in range(n_partitions)]

    def produce(self, key, message):
        """Route message to partition based on key."""
        p = partition_for(key, len(self.partitions))
        self.partitions[p].append(message)

    def read_partition(self, partition_id):
        return self.partitions[partition_id]

topic = FakeKafkaTopic(n_partitions=3)
topic.produce("ORD-1", {"order_id": "ORD-1", "event": "created"})
topic.produce("ORD-1", {"order_id": "ORD-1", "event": "paid"})
topic.produce("ORD-2", {"order_id": "ORD-2", "event": "created"})

# All ORD-1 events are in partition 2, in order
print("Partition 2:", topic.read_partition(2))
# [{"order_id": "ORD-1", "event": "created"}, {"order_id": "ORD-1", "event": "paid"}]

Sum for ORD-1: 79+82+68+45+49 = 323. 323 % 3 = 2.

Characterord() value
O79
R82
D68
-45
149

Common beginner questions

Why not use Python's built-in hash()?

CPython randomizes hash() across processes (for security). This means hash('ORD-1') returns a different number each time you restart Python. Partition assignment must be deterministic, so we use a stable hash like sum of ordinals or a proper hash function like murmur3.

How many partitions should a topic have?

More partitions means more parallelism but also more overhead. A common starting point is 6-12 partitions for a moderately busy topic. You can add partitions later, but you cannot easily reduce them.

What about ordering across partitions?

Kafka only guarantees ordering within a single partition. If you need ordering across all messages regardless of key, you need a single partition (which sacrifices parallelism).

hash() will give different results each run

Python's built-in hash() is randomized per process. The grader expects partition_for('ORD-1', 3) to return 2 every time. Use sum(ord(c) for c in key) instead.

Parallelism from the PySpark track

PySpark partitions data across worker nodes for parallel processing. Kafka partitions serve the same purpose: splitting work so multiple consumers can process in parallel.

What comes next

Now that you know where messages live (partitions), the next lesson explains how consumers track their position: offsets. An offset is a bookmark that tells the consumer which message to read next.

Practice

Implement partition_for, call it with 'ORD-1' and 3 partitions. Store 2 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?
Why a queueOffsets and commits