Partitioning Strategy — Default Behavior & Building a Custom Partitioner
What Happened Last Time
Look back at the delivery reports from your EventTrackingService script.
Every time user-501 sent an event, it landed on the same partition number. You never chose that partition. Kafka's default behavior decided it for you.
Today we open that decision up completely. Then we take control of it ourselves.
How the Default Partitioner Works
How does the producer choose a partition when a message has a key?
It runs the key through a hash function. A hash function is a piece of math that turns a key into a number, and it gives the same number every time for the same key. That number is then mapped onto one of the topic's partitions.
partition = hash(key) % number_of_partitions
Same key in, same partition out. That is why all events of one user stayed together.
Which hash function does confluent-kafka use?
CRC32. (Its default partitioner is called consistent_random.) You can predict the result yourself, because Python has the same function:
pythonimport zlib for u in range(501, 508): print(u, zlib.crc32(f"user-{u}".encode()) % 3)
For our seven users and a topic with three partitions, it prints this:
| User | Default partition (CRC32 % 3) |
|---|---|
| user-501 | 1 |
| user-502 | 0 |
| user-503 | 0 |
| user-504 | 1 |
| user-505 | 1 |
| user-506 | 1 |
| user-507 | 2 |
Compare this table with the delivery reports from the last lecture. They match.
What about messages that have no key?
There is nothing to hash. The producer then picks partitions on its own, and spreads the messages across the topic over time.
There are two things to remember about this default.
1. The hash is consistent, but it is opaque. It is always the same for a key, but it is not meaningful to you. You cannot look at it and say "I want this key on partition 0". It lands wherever the math says.
2. Not every Kafka client uses the same hash function. Some other client libraries use a hash called Murmur2, not CRC32. So the same key can land on a different partition depending on which client wrote it. If Python and another kind of producer write to the same topic, and something downstream assumes that "same key means same partition", you can get a quiet bug. The fix is a setting. confluent-kafka can use the Murmur2 hash too, which is what several other clients use by default:
pythonproducer_config = { "bootstrap.servers": BOOTSTRAP_SERVERS, "partitioner": "murmur2_random", }
(This is only an example. It is not part of the scripts in this lecture.)
And remember: the mapping holds only while the number of partitions stays the same. If you add partitions to a topic later, hash % number_of_partitions gives different answers, and keys move to new partitions.
Today's Requirement
Your platform just launched a real-time personalization model. It updates a shopper's recommendations the instant they click, not overnight.
Product wants to roll it out to premium subscribers first. The plan is a dedicated communication channel:
Output / NoteEvery clickstream event from a premium-tier subscriber must always go to partition 0. A dedicated, low-latency consumer will read only that partition and feed the personalization model. All other users' events keep spreading over the remaining partitions, read by the shared consumer.
Can the default partitioner do this?
No. It has no idea what "premium" means. And you cannot force one key onto one partition.
Default hashing vs. custom partitioning logic
Let's see how bad it is. Look at the table again. With three partitions and default hashing:
user-505is a premium user, but goes to partition 1. Wrong.user-503is a standard user, but goes to partition 0. Wrong.
Only user-502 is on partition 0 by luck. The default partitioner cannot meet our requirement.
The Real Mechanism
How do we take control?
In confluent-kafka, you cannot plug your own partitioner class into the producer. That is simply not how this client works. Instead, you decide the partition yourself, in plain Python, and tell produce() where to put the message.
pythonproducer.produce( topic=TOPIC, key=user_id.encode("utf-8"), value=json.dumps(event).encode("utf-8"), partition= partition_no, callback=delivery_report, )
That is all. Just one extra argument, partition=, on a call you already know.
Does this switch the default partitioner off for the whole producer?
No. It is a per-message choice. If you pass partition=, the default partitioner is skipped for that message. If you leave it out, the default partitioner works as usual. You could even mix both in one script.
Building It
We start from the same producer as before: same config, same callback, same poll() and flush(). We add three things: a way to know how many partitions exist, a function that decides the partition, and a tier for each user.
Step 1: Ask the Cluster How Many Partitions Exist
To reserve partition 0 and spread everyone else over "the rest", we need to know how many partitions there are. Writing 3 in the code is tempting. But if someone adds partitions to the topic later, that number is silently wrong. So we ask the cluster:
pythondef get_partition_count(producer: Producer, topic: str) -> int: cluster_metadata = producer.list_topics(topic, timeout=10) topic_metadata = cluster_metadata.topics[topic] if topic_metadata.error is not None: raise RuntimeError(f"Could not fetch metadata for '{topic}': {topic_metadata.error}") return len(topic_metadata.partitions)
Where does list_topics() come from?
It is a method of the Producer you already have (and of the Consumer). Every client keeps a connection to the cluster's metadata. With a topic name, it asks only about that topic, so you get its real, current partition list. The timeout=10 means: if the cluster does not answer in 10 seconds, stop waiting. And if the topic has an error, we raise a clear exception right away, instead of failing later with a confusing crash.
Step 2: Decide the Partition
pythondef decide_partition(user_id: str, tier: str) -> int: if tier == 'premium': return 0 remaining_partitions = get_partition_count(producer, TOPIC) - 1 partition_no = zlib.crc32(user_id.encode("utf-8")) % remaining_partitions return partition_no + 1
Let's read it slowly.
- Premium user: return
0. Always. - Standard user: we must pick partition 1 or 2, but always the same one for the same user. So:
remaining_partitionsis the number of partitions, minus partition 0. With three partitions, it is 2.zlib.crc32(...) % remaining_partitionsgives 0 or 1. It is a hash, so the same user gives the same answer every time.+ 1shifts it into partitions 1 and 2, and skips partition 0.
Why do we compute the partition for every message, and not only for premium ones?
Because the default partitioner cannot express our rule. It hashes over all partitions, including partition 0. If we let it handle standard users, some of them would land on partition 0 by chance. That would ruin the reserved partition.
Why zlib.crc32?
It is the same hash function that confluent-kafka uses by default. So our standard users keep the same "same key, same partition" property. We did not remove hashing. We took over the choice, so that partition 0 is excluded.
Output / NoteNote for real systems. The file asks the cluster for the partition count every time a standard message is sent. For a demo with 30 messages, that is fine. But it is a network request each time. In real code, read the count once and keep it. Read it again only from time to time:
pythonNUM_PARTITIONS = get_partition_count(producer, TOPIC) # once, at startupAlso, this logic needs at least 2 partitions. With only one,
remaining_partitionswould be 0, and the%would fail with a division by zero.
Step 3: Give Users a Tier
Our events now carry a tier, so decide_partition() has something to check:
pythondef make_click_event(user_id: str, tier: str) -> dict: return { "user_id": user_id, "tier": tier, "event_type": random.choice(EVENT_TYPES), "event_time": time.time(), }
And here are our users, with their tiers:
pythonusers = [ ("user-501", "standard"), ("user-502", "premium"), ("user-503", "standard"), ("user-504", "standard"), ("user-505", "premium"), ("user-506", "standard"), ("user-507", "standard") ]
Step 4: Put It Together
pythonfor _ in range(30): user_id, tier = random.choice(users) event = make_click_event(user_id, tier) partition_no = decide_partition(user_id, tier)
Then comes the produce() call with partition= partition_no, which you saw above. Everything else is unchanged: the config, the serialization, the callback, poll() and flush(). Custom partitioning is not a different way of producing. It is one extra decision, on top of what you know.
Running It
Step 1: Start Your Cluster
bashdocker compose start docker ps
The clickstream-events topic already exists from the last lecture. You do not need to create it again.
Step 2: Run the Script
In your 01_producers folder:
bashpython 2-premium_clickstream_producer.py
It sends 30 events, one per second. Watch the delivery reports.
What should you see?
With three partitions, the delivery reports must match this table:
| User | Tier | Partition |
|---|---|---|
| user-502 | premium | 0 |
| user-505 | premium | 0 |
| user-501 | standard | 1 |
| user-503 | standard | 1 |
| user-504 | standard | 2 |
| user-506 | standard | 2 |
| user-507 | standard | 2 |
Both premium users go to partition 0, every time. No standard user ever goes to partition 0. And each standard user goes to the same partition every time. The last part is the rule that Kafka's default partitioner gave us, and we kept it.
Step 3: Verify From the Consumer Side
bashdocker exec -it kafka-1 bash export PATH=$PATH:/opt/kafka/bin
bashkafka-console-consumer.sh \ --topic clickstream-events \ --bootstrap-server kafka-1:29092 \ --from-beginning \ --timeout-ms 5000 \ --property print.key=true \ --property print.partition=true \ --property key.separator=" | " | tail -30
Why tail -30?
The topic also holds the 20 events from the last lecture, which used the default hashing (there, user-505 went to partition 1). We only want the 30 newest events, the ones from this run. --timeout-ms 5000 makes the consumer stop by itself after 5 seconds without new messages.
Look for user-502 and user-505. In these 30 lines, every one of them shows partition 0.
Step 4: Stop Your Cluster
bashdocker compose stop
Things to Keep in Mind in Real Systems
Reserving a partition is a real, common pattern. But it comes with trade-offs.
- The rule lives in your code. Every producer that writes to this topic must follow the same rule. A producer that uses the default partitioner can put a standard user's event on partition 0.
- If a user changes tier, their events start going to a different partition. Ordering across that change is no longer guaranteed.
- A reserved partition is a skewed partition. Partition 0 may carry much less (or much more) traffic than the others.
- The dedicated consumer that reads only partition 0 must be told to read that partition explicitly. A normal group subscription would not do that. You will see how in the consumer chapter.
- Adding partitions changes the mapping of standard users, just as it does for the default hashing.
Checklist
- Understood that default partitioning is
hash(key) % number_of_partitions: consistent, but opaque - Know that
confluent-kafkauses CRC32, and that other clients may use a different hash (and how to switch withpartitioner) - Understood why the default partitioner cannot meet the "premium on partition 0" requirement
- Understood that in this client you choose the partition yourself and pass
partition=toproduce(), per message - Wrote
get_partition_count()anddecide_partition(), and understood every line - Understood why we hash standard users ourselves, with the same function as the default
- Ran the script, and confirmed premium events always land on partition 0
- Know the real-world trade-offs of reserving a partition
You just steered data placement on purpose. This is a real pattern for tiered, low-latency pipelines, and not just Kafka trivia. Next: we look at what happens after a message reaches the broker. We study delivery guarantees, and how the idempotent producer prevents duplicate writes.