Topics, Partitions & Keys
Ordering is per partition, not per topic. Keys decide partition assignment, which decides ordering locality. Learn this until it's instinct.
SPACED REPETITION Β· 15 practice questions
Make this lesson stick.
Try 3 questions now. No account needed. Sample answers aren't saved.
or sign in to practice all 15By the end of this lesson you will be able to:
- explain why ordering, storage and parallelism belong to a partition rather than a topic;
- predict which partition a keyed record lands on, and why a .NET and a Java producer can disagree;
- say what the default .NET producer does with records that have no key;
- size a topic's partition count and replication factor, and create and inspect it from C#;
- work out what breaks β and for which keys β when you add partitions to a live topic.
Prerequisite: "Core Kafka Mental Model" (the log, offsets, retention and replication basics). Most major sections include a short exercise β try it before opening the answer.
Introduction: The Topic-Partition-Key Mental Model
If you've built messaging into a .NET system before, you already have a mental model for what a "queue" does: a producer drops a message in, a consumer picks it up, and once it's processed, it's gone. MSMQ works this way. So does RabbitMQ. So does Azure Service Bus. Kafka calls its named streams "topics," which sounds queue-adjacent, but the internal structure is genuinely different, and that difference has consequences from the very first line of producer code.
If you've completed "Core Kafka Mental Model," the next subsection ("Kafka Topics Are Logs, Not Queues") will be familiar β skim it and pick up at "Why the Partition, Not the Topic, Is the Real Unit." From there on the material is new: the physical anatomy of a topic, the exact routing rules the .NET client uses (which are not the ones Java-oriented material describes), and the operational traps that partition and key decisions create.
Three questions frame the lesson. What actually gets stored when you publish a message? What decides which part of the system a given message ends up on? And why should a .NET developer care about the distinction? The answers hinge on three words β topic, partition, and key.
Kafka Topics Are Logs, Not Queues
A Kafka topic is not a container that holds messages until they're consumed. It's a named, append-only commit log β records are written to the end and stay there, in arrival order, identified by a numeric position called an offset. A topic isn't one log file, though: it's split into one or more partitions, each its own independent, ordered, append-only log, spread across the brokers in your cluster.
The consequence that matters most here is retention. In MSMQ, RabbitMQ, or Azure Service Bus, a message is transient β it exists to be delivered, and acknowledgment (or expiry) removes it. Kafka inverts this. A record stays at its offset for as long as the topic's retention policy allows, independent of whether zero consumers or a hundred have read it. A consumer doesn't take a message; it reads a copy and advances its own tracked offset. Another consumer group can read the same records from the beginning later.
π‘ Mental Model: A queue is a checkout line β once you're served, you leave and the line forgets you. A Kafka partition is a numbered ledger β every entry stays on the page at its line number, and "reading" means glancing at line 47 without erasing it. When that ledger gets trimmed is a retention-policy question (see "Core Kafka Mental Model"), not a property of the log itself.
Why the Partition, Not the Topic, Is the Real Unit
Here's the detail that trips up newcomers: when people describe a topic's behavior β its ordering, parallelism, throughput β they're almost always describing partition behavior. The topic is a logical name you subscribe to; nearly everything operationally significant happens one level down:
π§ Storage β each partition is a physically separate log on disk (replicated across brokers), not a shared structure.
π Ordering β offsets are assigned per partition, so "record 10 comes after record 9" holds within a single partition, never across the topic. (Its limits β producer retries that can reorder, and what a consumer sees across a rebalance, when the group reshuffles which consumer owns which partition β are covered in "Core Kafka Mental Model" and "Consumer Groups & Offsets"; internalize for now that ordering is partition-scoped.)
π§ Parallelism β a consumer within a group is assigned whole partitions, so partition count caps how many consumers in that group can work simultaneously. (This is the rule for ordinary consumer groups, which is what GroupId gives you. Kafka 4.2's share groups let several consumers share one partition, at the cost of per-partition ordering; they are a different tool, and one your Confluent.Kafka version may not support yet, so check before planning around them.)
π― Routing β every record lands in exactly one partition, and the mechanism deciding which is the partitioner, driven primarily by the record's key.
That last point is the hinge. If two records carry the same key, the partitioner routes them to the same partition β the same ordered log β so their relative order is preserved (provided the producer itself does not reorder them on retry; "Core Kafka Mental Model" covers the idempotence setting that prevents that). Different keys (or no key) give no such guarantee. So "will my messages be processed in order?" is really "did my messages land on the same partition?", which is a direct function of key choice and the partitioning mechanism dissected in "How the Partitioner Assigns Keys to Partitions."
π― Key Principle: Kafka does not order a topic. It orders each partition independently. Any correctness argument that depends on message ordering has to be phrased in terms of "same partition," not "same topic."
A Running Example: Order Events
Every code sample in this lesson builds against one scenario: an e-commerce system publishing order events β OrderCreated, OrderPaid, OrderShipped β to a topic named order-events, keyed by order ID. That choice isn't arbitrary; it makes the partitioner observable, because events sharing a key land on the same partition and a consumer reading that partition sees them in production order. Whether OrderId is the best key (versus, say, CustomerId) is a design question taken up in "Choosing a Key: Three Questions" near the end of this lesson. Until then it exists purely to make routing tangible.
// The event payload our producer will serialize as the message value.
// The OrderId will separately be used as the record's Kafka key.
public record OrderEvent(
string OrderId,
string EventType, // e.g. "Created", "Paid", "Shipped"
DateTimeOffset Timestamp,
decimal Amount
);
And a preview of a produce call using Confluent.Kafka's IProducer<TKey, TValue> β full setup and delivery-result handling come later, but seeing the key/value split now makes the relationship concrete:
// Illustrative shape only β producer construction and config are covered later.
// Note that Key and Value are separate fields: the key drives partition routing,
// the value is the payload a consumer will deserialize.
var message = new Message<string, string>
{
Key = orderEvent.OrderId, // routing input
Value = JsonSerializer.Serialize(orderEvent) // payload
};
var deliveryResult = await producer.ProduceAsync("order-events", message);
// deliveryResult.Partition tells you which partition the partitioner chose
Key and Value are distinct fields β a distinction that doesn't exist in the same form in queue-based .NET messaging APIs, where a body is just a body. In Kafka the key is not incidental metadata; it's the input to the routing decision that determines which ordered log the record joins.
What This Lesson Covers β and What It Deliberately Doesn't
This lesson walks through the physical anatomy of a topic in "Anatomy of a Topic: Partitions, Brokers, and Replicas"; the mechanics of turning a key into a partition number β and the places where the .NET client's behavior diverges from what Java-oriented documentation describes β in "How the Partitioner Assigns Keys to Partitions"; full producer/consumer code in "Producing and Consuming with Confluent.Kafka: A Worked Example"; and the operational mistakes these decisions cause in "Common Pitfalls: Partition Count, Skew, and Remapping."
Two closely related topics are deliberately kept light here: the exact limits of per-partition ordering (producer retries are covered in "Core Kafka Mental Model", rebalances in "Consumer Groups & Offsets") and how long records persist before deletion or compaction ("Core Kafka Mental Model"). Choosing which field to key by gets a short practical section at the end. Keeping the rest separate lets this lesson stay on one job: building the structural model β topic, partition, broker, key β that every later lesson assumes.
π‘ Real-World Example: Picture a PaymentProcessed event and a ShipmentDispatched event for two different orders, produced back to back. Because they carry different order IDs as keys, the partitioner may route them to different partitions entirely β and a consumer has no guarantee it will see them in production order, because they aren't competing for position in the same log. Reasoning about "did A happen before B" in Kafka always starts with "were A and B on the same partition?"
β οΈ Common Mistake: Treating a Kafka topic as if it behaves like an Azure Service Bus topic with subscriptions, where the middleware quietly manages delivery and cleanup for you. β Wrong thinking: "I don't need to think about partitions β Kafka will just deliver each message to my consumer and clean up after itself." β Correct thinking: "My topic is a set of independently ordered logs; my consumer is assigned specific partitions to read, my messages persist until a retention policy removes them, and which partition each message lands on is a routing decision I make (or delegate) through the key."
βοΈ Exercise: which pairs are ordered?
A producer awaits each of these sends in turn, to order-events, with the default partitioner and idempotence enabled. For each pair, is a consumer guaranteed to read the first record before the second?
OrderCreatedthenOrderPaid, both with keyORD-4471.OrderCreatedfor keyORD-4471, thenOrderCreatedfor keyORD-90110.- Two records with
Key = null, sent back to back. OrderCreatedforORD-4471from a .NET service, thenOrderPaidforORD-4471from a Java service writing to the same topic.
Check your answer
- Yes. Same key, same partition count, same partitioner β same partition, and the idempotent producer does not reorder on retry.
- No. Different keys may hash to different partitions, and across partitions there is no order. (Even if they happen to share a partition, relying on that is relying on a coincidence of the hash.)
- No. Unkeyed records are spread over partitions; two of them may or may not end up on the same one.
- No. The .NET client hashes keys with CRC32 and the Java client with murmur2, so the two services usually send the same key to different partitions β the subject of "Divergence 1" below.
Anatomy of a Topic: Partitions, Brokers, and Replicas
A topic looks deceptively simple from the outside β you name it, you send to it, you read from it. But a topic is not a single file or a single queue on one machine. It is a logical name mapping onto one or more partitions, and those partitions are the physical units that get stored, replicated, and read. Understanding this split β topic as label, partition as reality β is the prerequisite for everything else, because every guarantee Kafka makes, and every guarantee it deliberately does not make, is a statement about partitions.
A Partition Is an Ordered, Immutable Log
Picture each partition as an append-only file. New records are written to the end, existing records are never modified, and every record carries a sequential integer called an offset marking its position within that specific partition.
Topic: order-events (6 partitions β P0 to P2 shown)
Partition 0:
offset 0 β OrderCreated(id=1042)
offset 1 β OrderShipped(id=1042)
offset 2 β OrderCreated(id=1050)
Partition 1:
offset 0 β OrderCreated(id=1043)
offset 1 β OrderCancelled(id=1043)
Partition 2:
offset 0 β OrderCreated(id=1044)
offset 1 β OrderShipped(id=1044)
offset 2 β OrderCreated(id=1046)
offset 3 β OrderCreated(id=1051)
Notice that offset 1 exists in every partition shown and refers to three completely different records. Offsets are per-partition counters, not one sequence shared across the topic. There is no such thing as "offset 1 of order-events" β you must say "offset 1 of partition 1," because that's the only coordinate system Kafka maintains. A developer coming from a single-queue model (an MSMQ queue with one linear backlog) has to unlearn the idea that a topic has one running position counter; it has as many counters as it has partitions, all advancing independently.
π― Key Principle: A topic's ordering and storage guarantees exist at the partition level, not the topic level. "Topic" is an organizational label over a set of independently-ordered logs.
Leaders, Followers, and Replication
Partitions don't live on just one broker if you want fault tolerance β and in any real deployment, you do. For each partition, Kafka designates one broker as the leader, which handles all writes and, by default, all reads for that partition. The remaining copies live on other brokers as follower replicas, continuously fetching from the leader. (Consumers can be configured to read from a nearby follower instead β rack-aware fetching, which needs the consumer's client.rack, a broker.rack on each broker, and the broker-side replica.selector.class β but producers always write to the leader.) The number of copies β leader plus followers β is the topic's replication factor.
Partition 0 (replication factor = 3)
Broker 1: [Leader] β handles all reads/writes for Partition 0
Broker 2: [Follower] β replicates Partition 0's log
Broker 3: [Follower] β replicates Partition 0's log
If the broker hosting the leader fails, one of the in-sync followers is promoted, and clients transparently redirect. Replication factor 3 is the common production baseline because it lets you lose one broker and still have two copies of every record β which is what keeps a topic writable under the common production setting min.insync.replicas=2 (the broker default is 1). (Leader election picks from the in-sync replica set, and it is not a majority vote β a single surviving in-sync replica can become leader. So with replication factor 3, a two-broker outage still leaves the partition readable; what stops is acks=all writes, because the ISR has shrunk below min.insync.replicas=2. The mechanics of in-sync replica sets and acknowledgment settings are in "Core Kafka Mental Model".)
The structural point for this lesson is simpler: replication factor is set per topic and applies to every one of its partitions β it is how many broker-local copies of each partition's log exist. A topic with 6 partitions and replication factor 3 doesn't store 3 copies of the topic β it stores 3 copies of each of the 6 partitions, for 18 physical partition-copies scattered across the cluster.
β οΈ Common Mistake: Assuming replication factor and partition count are the same knob. Partition count answers "how many parallel logs make up this topic"; replication factor answers "how many broker-local copies does each of those logs get." Setting replication factor to 1 outside local experimentation means a single broker failure makes that partition's data unavailable β there's no follower to promote.
Partition Count Sets the Ceiling on Parallelism
Because each partition can only be actively consumed by one consumer within a given consumer group at a time (assignment mechanics are covered in the worked example; share groups, mentioned earlier, are the exception), partition count puts a hard ceiling on how many consumer instances in a single group can do useful work simultaneously. A topic with 3 partitions keeps at most 3 consumers in one group busy; a fourth sits idle, because there's no fourth partition to hand it. This is the structural reason partition count is a capacity-planning decision rather than a storage detail.
π‘ Mental Model: Think of partitions as a fixed number of assembly lines. Adding more workers than lines doesn't speed anything up β the extras stand around. The number of lines is decided when the factory is built, and while Kafka lets you add lines later, doing so has a real cost for keyed messages, as you'll see at the end of this lesson.
Creating a Topic from .NET Code
All of this structure has to be specified somewhere, and in a Confluent.Kafka application that's typically the IAdminClient interface rather than a broker's command-line tooling. CreateTopicsAsync declares partition count and replication factor programmatically, which is useful for integration tests, provisioning scripts, or any path where topic creation should be deployable code rather than a manual step.
using Confluent.Kafka;
using Confluent.Kafka.Admin;
var adminConfig = new AdminClientConfig
{
BootstrapServers = "localhost:9092"
};
using var admin = new AdminClientBuilder(adminConfig).Build();
try
{
// Read from environment-specific settings (see the warning below):
// 1 on a single-broker dev setup, 3 in production.
short replicationFactor = settings.ReplicationFactor;
await admin.CreateTopicsAsync(new[]
{
new TopicSpecification
{
Name = "order-events",
NumPartitions = 6, // parallelism ceiling for one consumer group
ReplicationFactor = replicationFactor // 3 = leader + 2 followers per partition
}
});
Console.WriteLine("Topic 'order-events' created.");
}
catch (CreateTopicsException ex)
{
// CreateTopicsAsync throws if the topic already exists or the
// requested replication factor exceeds the number of available brokers
Console.WriteLine($"Topic creation failed: {ex.Results[0].Error.Reason}");
}
Three things are worth calling out. NumPartitions and ReplicationFactor are the two structural properties this section has been building toward β everything else in TopicSpecification (retention and compaction overrides) is orthogonal to anatomy. CreateTopicsAsync is not idempotent: calling it for a topic that already exists throws a CreateTopicsException rather than silently succeeding, so provisioning code typically checks existing topics first or catches and inspects that exception. And requesting a replication factor higher than the number of brokers currently in the cluster fails outright β you cannot replicate onto more brokers than exist, which is a concrete illustration of replication factor as a physical constraint rather than a configuration number.
β οΈ Common Mistake: Hardcoding ReplicationFactor = 1 because that's what local single-broker Docker setups need, then deploying the same code unmodified against a multi-broker cluster. The topic still gets created, but it silently forfeits the fault tolerance the cluster exists to provide. Make replication factor a configuration value read from environment-specific settings, precisely so a local dev value can't leak into a production provisioning script.
Reading the Anatomy Back Out
Once a topic exists, the same admin client can describe its structure β confirming that the partition and replica layout matches what you requested, and revealing the current leader for a given partition. That workflow, using GetMetadata(), appears in the worked example later, when you'll have real data to inspect. For now the model to carry forward is architectural: a topic is a named collection of independently-ordered partitions, each with exactly one leader broker and zero or more follower replicas, and partition count is the hard limit on how much of that topic a single consumer group can process in parallel. Keys β which decide which partition a record lands in β are the mechanism connecting this physical layout to your application's data.
βοΈ Exercise: size the topic
order-events is created with 6 partitions and replication factor 3 on a 3-broker cluster, with min.insync.replicas=2.
- How many partition copies does the cluster store for this topic, and how many does each broker hold if they are spread evenly?
- The consumer group runs 8 instances. How many are doing work?
- Two brokers go down at once. Can consumers still read? Can an
acks=allproducer still write? - Harder: traffic is expected to triple next year and each consumer instance can handle about one partition's worth of peak traffic today. What would you change now, and why now rather than later?
Check your answer
- 6 Γ 3 = 18 partition copies; 6 per broker β with 3 brokers and RF 3, every broker holds a copy of every partition (as leader for about 2 of them and follower for the rest).
- 6. Two instances sit idle: there are only 6 partitions to hand out.
- Reads: yes β the one surviving in-sync replica of each partition becomes (or stays) leader.
acks=allwrites: no β the ISR is down to 1, belowmin.insync.replicas=2, so the leader rejects them withNOT_ENOUGH_REPLICAS. librdkafka retries that error quietly, soProduceAsyncjust waits: it succeeds if a broker is back in sync withinmessage.timeout.ms(5 minutes by default), and otherwise throws aProduceExceptionwithLocal_MsgTimedOut. - Create the topic with more partitions from the start (for example 18 or 24) rather than planning to add them later: adding partitions to a keyed topic remaps keys (see "What Breaks When You Add Partitions"), and the count can never be reduced. Keep the number modest β every partition is standing broker overhead.
How the Partitioner Assigns Keys to Partitions
When a .NET producer calls ProduceAsync on a topic with several partitions, something decides which partition receives the record. That component is the partitioner, and understanding it is what lets you predict β rather than guess β where a message ends up. This section is deliberately narrow: it covers the mechanics of routing, not which field makes a good key (that comes in "Choosing a Key: Three Questions").
It is also the section where Confluent.Kafka diverges most sharply from the Java client, in two ways that silently break systems. Read the divergence subsections even if you already know how Kafka partitioning works "in general."
The Default Partitioner: Hashing a Key
When you produce a message with a non-null key, the default partitioner applies a hash-modulo formula: it hashes the key's bytes and takes that hash modulo the current number of partitions, yielding a partition number between 0 and partitionCount - 1.
key bytes β hash function β hash value
hash value % partitionCount β partition number
The critical property is determinism: the same key, hashed against the same partition count, always produces the same partition number. With librdkafka's default partitioner, OrderId = "ORD-4471" lands on partition 3 of the 6-partition order-events topic today, and on partition 3 tomorrow and every time after β as long as the topic's partition count doesn't change. This is what makes keyed ordering possible: every producer instance using the same partitioner, running independently, computes the identical mapping, so all records for that order converge on one partition without any coordination between producers.
using Confluent.Kafka;
var config = new ProducerConfig { BootstrapServers = "localhost:9092" };
using var producer = new ProducerBuilder<string, string>(config).Build();
// The key ("ORD-4471") is hashed by the default partitioner.
// Every message with this exact key lands on the same partition.
var result = await producer.ProduceAsync("order-events", new Message<string, string>
{
Key = "ORD-4471",
Value = "{ \"status\": \"Created\" }"
});
Console.WriteLine($"Delivered to partition {result.Partition.Value} at offset {result.Offset.Value}");
Run this three more times with the same key and different payloads ("Paid", "Shipped", "Delivered") and every one reports the same partition through result.Partition. The value being produced is irrelevant to routing; only the key bytes matter.
β οΈ Common Mistake: Assuming the partition number is stable across topic resizes. The modulo depends on the current partition count, so changing that count changes the formula's output for existing keys. The operational consequences are covered in "Common Pitfalls" later; the mechanism to remember is that the formula is hash(key) % partitionCount, not hash(key) % (some fixed number).
Divergence 1: Which Hash Function .NET Actually Uses
Here is the fact most Kafka material will not tell you, because most Kafka material is written for the JVM.
Confluent.Kafka wraps librdkafka, whose default partitioner setting is consistent_random β a CRC32-based hash. The Java producer uses murmur2. These are different functions, so the same key produces a different partition number depending on which client produced it.
That is fine in a homogeneous .NET shop and catastrophic in a mixed one. If a .NET service and a Java (or Kafka Streams, or Spark, or Flink) service both produce ORD-4471 events to the 6-partition topic, the .NET events land on partition 3 and the Java events on partition 4, and the per-key ordering guarantee you thought you had never existed. (For any particular key the two hashes agree only by coincidence β with 6 partitions, roughly one key in six β so "usually different" is the right expectation.) Nothing errors. Nothing logs. A consumer just sees an order's history split across two logs with no relative order between them.
librdkafka exposes the alternatives through ProducerConfig.Partitioner:
βοΈ Partitioner value | π Keyed routing | π³οΈ Null key goes to |
|---|---|---|
ConsistentRandom (default) | CRC32 (an empty key is treated like a null key) | A sticky random partition, changed every few ms |
Consistent | CRC32 | One fixed partition |
Murmur2Random | murmur2 β Java-compatible | A sticky random partition, changed every few ms |
Murmur2 | murmur2 | One fixed partition |
Random | A random partition per record | A sticky random partition, changed every few ms |
// Any topic that is ALSO produced to by a JVM client must use murmur2,
// or the two producers will disagree about where a key belongs.
var config = new ProducerConfig
{
BootstrapServers = "localhost:9092",
Partitioner = Partitioner.Murmur2Random
};
π― Key Principle: Cross-client key colocation is opt-in, not automatic. If a topic has producers on more than one client stack, pin the partitioner explicitly on every one of them and write down which you chose. Murmur2Random is the setting that matches the Java default's keyed routing.
β οΈ Common Mistake: Discovering this after the fact and "fixing" it by switching the .NET producer to Murmur2Random on a live topic. That has the same effect as a partition-count change β most existing keys' target partitions move (all but the ones where the two hashes happen to agree), and records already written stay where they are. Treat it as a migration, not a config tweak. On a new topic, choosing it up front costs nothing.
Divergence 2: What Happens When the Key Is Null
Not every message needs a key. With Key = null there's no hash to compute, so the partitioner falls back to a different strategy β and while both clients use the same idea, they decide when to switch differently.
Both use sticky partitioning: rather than picking a new partition for every record, the producer sends null-key records to one partition for a while, then switches. The reasoning is throughput β producers batch per partition, and rotating on every record spreads each partition's queue thin, producing many small batches.
- librdkafka (Confluent.Kafka) picks a random available partition and sticks to it for
sticky.partitioning.linger.ms, which defaults to twicelinger.msβ 10 ms with the 5 ms default linger. Then it picks another. Setting it to0turns stickiness off and gives a random partition per record. - The Java client (KIP-480, refined in Kafka 3.3 by KIP-794; KIPs are Kafka Improvement Proposals) switches after about
batch.sizebytes (the producer's maximum batch size) have gone to the current partition, not after a time interval.
librdkafka default (consistent_random), null key, steady traffic:
records produced in ms 0β10 β P3
records produced in ms 10β20 β P0
records produced in ms 20β30 β P3 (random choice each switch)
Java client (sticky, KIP-794), null key:
first ~batch.size bytes β P1
next ~batch.size bytes β P4
Client defaults do shift between versions, so check sticky.partitioning.linger.ms and partitioner against the librdkafka version your package pins if this matters to your throughput budget. What does not shift is the guarantee: if you don't supply a key, you get reasonable load distribution and no ordering relationship whatsoever between any two unkeyed messages, since they may not even share a partition. That falls directly out of the routing mechanism; it isn't a separate rule.
π‘ Mental Model: Think of the default partitioner as a two-mode switch. Key present β deterministic CRC32 routing. Key absent (or empty) β a randomly chosen partition that changes every few milliseconds. Nothing more exotic happens by default.
βοΈ Exercise: predict the partitions
Here are the real CRC32 values librdkafka computes for four keys (unsigned 32-bit):
| Key | CRC32 |
|---|---|
ORD-4471 |
297785037 |
ORD-48213 |
3289117203 |
ORD-90110 |
1650520494 |
ORD-10233 |
4217976590 |
- With the default partitioner and 6 partitions, which partition does each key go to?
- Which two keys share a partition, and does that mean their records are ordered relative to each other?
- A burst of 10,000 null-key records is produced within 25 ms. Roughly how are they spread across partitions with the default settings?
Check your answer
- Partition = CRC32 % 6:
ORD-4471β 3,ORD-48213β 3,ORD-90110β 0,ORD-10233β 2. ORD-4471andORD-48213both land on partition 3, so yes, a consumer of partition 3 reads them in write order. But that is a coincidence of the hash, not a guarantee you designed β a partition-count change or a different partitioner could separate them.- Mostly in a few large runs: the sticky partition changes about every 10 ms, so the burst lands on about three partitions (a new random pick at each switch, possibly the same one twice) rather than being sprinkled evenly over all six. Over longer periods the switching evens out.
Bypassing the Partitioner: Explicit Partition Targeting
Sometimes you already know which partition a record belongs on. Confluent.Kafka supports this through a ProduceAsync overload taking a TopicPartition instead of a topic name. The hashing partitioner is never invoked; the client sends the record straight to the partition you named.
using Confluent.Kafka;
var config = new ProducerConfig { BootstrapServers = "localhost:9092" };
using var producer = new ProducerBuilder<string, string>(config).Build();
// TopicPartition targets partition 1 directly.
// The key is still stored with the message, but it plays no role
// in routing here β the partitioner is bypassed entirely.
var topicPartition = new TopicPartition("order-events", new Partition(1));
var result = await producer.ProduceAsync(topicPartition, new Message<string, string>
{
Key = "ORD-4471",
Value = "{ \"status\": \"Refunded\" }"
});
Console.WriteLine($"Explicitly written to partition {result.Partition.Value}");
This is useful for migrating data between topics while preserving each record's original partition, or for diagnostic tooling that writes a test record to a specific partition. It's a narrow escape hatch rather than a routine pattern β manual assignment makes you responsible for whatever colocation guarantees your keys were supposed to provide.
β οΈ Common Mistake: Mixing explicit TopicPartition targeting with keyed produces for the same logical key. If one code path routes "ORD-4471" through the hash (landing on partition 3 of the 6-partition topic) while another forces it onto partition 1, you've silently broken the colocation guarantee keying was meant to provide β messages for one order now live on two partitions with two independent offset sequences.
Writing a Custom Partitioner
The default covers the overwhelming majority of cases, but you can override it entirely. Confluent.Kafka does this with a delegate, not an interface β ProducerBuilder.SetDefaultPartitioner (all topics) or SetPartitioner (one named topic), both taking a PartitionerDelegate:
Partition PartitionerDelegate(
string topic, int partitionCount, ReadOnlySpan<byte> keyData, bool keyIsNull);
Here is one that reserves partition 0 for VIP traffic and hashes everything else across the rest:
using Confluent.Kafka;
// FNV-1a: deterministic across processes, machines, and restarts.
static uint Fnv1a(ReadOnlySpan<byte> data)
{
const uint offsetBasis = 2166136261, prime = 16777619;
uint hash = offsetBasis;
foreach (var b in data)
{
hash ^= b;
hash *= prime;
}
return hash;
}
static Partition VipAwarePartitioner(
string topic, int partitionCount, ReadOnlySpan<byte> keyData, bool keyIsNull)
{
// Single-partition topics have nothing to reserve against, and the
// modulo below would divide by zero.
if (partitionCount < 2) return new Partition(0);
// Spread unkeyed records over the non-reserved partitions.
// Only reached if sticky partitioning is off (see the config below).
if (keyIsNull) return new Partition(Random.Shared.Next(1, partitionCount));
// UTF-8 literal β compares bytes without allocating a string.
if (keyData.StartsWith("VIP-"u8)) return new Partition(0);
return new Partition(1 + (int)(Fnv1a(keyData) % (uint)(partitionCount - 1)));
}
var config = new ProducerConfig
{
BootstrapServers = "localhost:9092",
// Without this, librdkafka never calls a custom partitioner for null keys:
// it routes them with its own sticky random choice over ALL partitions,
// including the reserved partition 0.
StickyPartitioningLingerMs = 0
};
using var producer = new ProducerBuilder<string, string>(config)
.SetDefaultPartitioner(VipAwarePartitioner)
// ...or scope it to one topic:
// .SetPartitioner("order-events", VipAwarePartitioner)
.Build();
Note the StickyPartitioningLingerMs = 0 line. librdkafka handles null-key records before your delegate: while sticky partitioning is on (the default), it picks a sticky random partition itself and your keyIsNull branch never runs. That is easy to miss because nothing fails β unkeyed traffic simply appears on the partition you meant to reserve.
β οΈ Common Mistake: Turning the key bytes back into a string and calling GetHashCode() on it inside a custom partitioner. .NET randomizes string hash codes per process by default, so two instances of the same service β or the same service before and after a restart β will route an identical key to different partitions. The bug is invisible in a single-process test and destroys colocation the moment you scale out. A custom partitioner must use a hash that is stable across processes: FNV-1a as above, CRC32, or murmur2.
Custom partitioners are a specialized tool. Reach for one only when you have a routing requirement the hash formula can't express, and remember it inherits all the same colocation responsibilities β a hand-rolled partitioner can scatter one order's events across partitions just as easily as a misused built-in.
π― Key Principle: The partitioner answers "which partition?" It has no opinion on "which field should be the key?" That second question β OrderId versus CustomerId and the trade-offs between them β is taken up in "Choosing a Key: Three Questions" at the end of the lesson. Here, treat the key as a given input and trace how it flows through hashing, random placement, explicit targeting, or custom logic to produce a partition number.
Producing and Consuming with Confluent.Kafka: A Worked Example
Everything so far is invisible until you run a producer and a consumer and watch the numbers come back. This section builds that end-to-end picture with a single scenario: order events keyed by OrderId. By the end you'll see exactly where partition and offset surface in the API, and confirm with real output that every event for an order lands on the same partition.
Building the Producer
The entry point is ProducerBuilder<TKey, TValue>. You configure broker addresses, specify key and value types as generic parameters, and Confluent.Kafka handles serialization with built-in serializers (Serializers.Utf8 for strings) or your own.
using Confluent.Kafka;
var producerConfig = new ProducerConfig
{
BootstrapServers = "localhost:9092"
};
using var producer = new ProducerBuilder<string, string>(producerConfig)
.SetKeySerializer(Serializers.Utf8)
.SetValueSerializer(Serializers.Utf8)
.Build();
// An order event: key = OrderId, value = a JSON payload (simplified as a string here)
var orderId = "ORD-48213";
var message = new Message<string, string>
{
Key = orderId,
Value = "{\"orderId\":\"ORD-48213\",\"status\":\"Created\"}"
};
DeliveryResult<string, string> result =
await producer.ProduceAsync("order-events", message);
Console.WriteLine(
$"Delivered to partition {result.Partition.Value} " +
$"at offset {result.Offset.Value}");
ProduceAsync sends the keyed message and awaits the broker's acknowledgment. The returned DeliveryResult<string, string> is where routing becomes observable: result.Partition is the partition the default partitioner chose (recall hash(key) % partitionCount), and result.Offset is the exact position the record now occupies in that partition's log. Call it again with the same orderId and you get the same partition number back every time β determinism as a concrete integer in your console output.
β οΈ Common Mistake: Using the fire-and-forget Produce method (which takes a delivery-report callback instead of returning a Task) and then letting the process exit before the callback fires. Buffered messages are lost silently. Note the trade-off, though: await-ing every ProduceAsync individually serializes your pipeline into one round trip per record and collapses throughput. Use ProduceAsync where a request path genuinely needs the result inline; use Produce with a delivery handler plus an explicit Flush(timeout) on shutdown for bulk paths.
Building the Consumer
ConsumerBuilder<TKey, TValue> mirrors the producer's shape. The critical extra configuration is GroupId, placing this consumer into a consumer group β a named set of consumers that split a topic's partitions between them.
using Confluent.Kafka;
var consumerConfig = new ConsumerConfig
{
BootstrapServers = "localhost:9092",
GroupId = "order-processing-service",
AutoOffsetReset = AutoOffsetReset.Earliest
};
using var consumer = new ConsumerBuilder<string, string>(consumerConfig)
.SetKeyDeserializer(Deserializers.Utf8)
.SetValueDeserializer(Deserializers.Utf8)
.Build();
consumer.Subscribe("order-events");
try
{
while (!cancellationToken.IsCancellationRequested)
{
ConsumeResult<string, string> record = consumer.Consume(cancellationToken);
Console.WriteLine(
$"key={record.Message.Key} " +
$"partition={record.Partition.Value} " +
$"offset={record.Offset.Value} " +
$"value={record.Message.Value}");
}
}
catch (OperationCanceledException) { /* graceful shutdown */ }
finally
{
consumer.Close(); // leaves the group immediately instead of waiting
// for the session timeout to expire
}
Each Consume call blocks until a record is available and returns a ConsumeResult<string, string>. The fields that matter for this model are record.Partition, record.Offset, and record.Message.Key. Note the symmetry: the partition and offset the consumer reports are exactly the values the producer's DeliveryResult reported when the record was written. Offsets are never reassigned or renumbered β they're a fixed coordinate in that partition's log.
Consumer Groups and Partition Assignment
A single consumer can read every partition of a topic, but production systems run several instances sharing a GroupId to parallelize. Kafka's rule is strict: within a consumer group, each partition is assigned to exactly one consumer at a time. With six partitions and two consumers, a typical split is three partitions each; scale to six consumers and each gets one; add a seventh and it sits idle, because a partition cannot be split between two members of a consumer group.
order-events (6 partitions)
Consumer group: order-processing-service (range assignor, 2 members)
Partition 0 ββ
Partition 1 ββΌβββΊ Consumer A
Partition 2 ββ
Partition 3 ββ
Partition 4 ββΌβββΊ Consumer B
Partition 5 ββ
Which consumer gets which partitions is decided by a pluggable partition assignor, configured via PartitionAssignmentStrategy. librdkafka defaults to range,roundrobin (the first strategy every member supports wins, so normally range; roundrobin deals partitions out to members in turn). The Range assignor hands out contiguous partition ranges per topic, which can produce mild imbalance when a group subscribes to several topics. The CooperativeSticky assignor minimizes partition movement during rebalances β when a member joins or leaves, it leaves as many existing assignments untouched as it can, rather than reshuffling everything, which shortens the processing pause a rebalance causes.
β οΈ Common Mistake: Switching an existing group to cooperative-sticky with a rolling deploy. librdkafka does not allow eager assignors (range, roundrobin: every member gives up all its partitions at each rebalance) and cooperative ones (cooperative-sticky: only the partitions that move are given up) in the same PartitionAssignmentStrategy list, so the two-step rolling migration described for the Java client is not available in Confluent.Kafka β and while old and new instances are mixed, the new ones fail to join the group. Plan it as a stop-all-then-start deploy (a short pause), or move to a new group ID with its offsets set deliberately. Newer clusters also offer the KIP-848 consumer protocol (GroupProtocol = GroupProtocol.Consumer), where the broker computes assignments and PartitionAssignmentStrategy no longer applies; "Consumer Groups & Offsets" covers it.
π‘ Mental Model: A consumer group is a team splitting a stack of numbered folders. Every folder is held by exactly one member at any moment β none shared, none skipped β but which member holds which can change as people join or leave.
Full Worked Example: Order Events Colocated by OrderId
Now put both sides together and confirm the claim that matters most: all events for the same order land on the same partition. order-events has its 6 partitions, and you produce three events for ORD-48213 and one for ORD-90110:
string[] orderIds = { "ORD-48213", "ORD-48213", "ORD-48213", "ORD-90110" };
string[] statuses = { "Created", "Paid", "Shipped", "Created" };
for (int i = 0; i < orderIds.Length; i++)
{
var msg = new Message<string, string>
{
Key = orderIds[i],
Value = $"{{\"orderId\":\"{orderIds[i]}\",\"status\":\"{statuses[i]}\"}}"
};
var result = await producer.ProduceAsync("order-events", msg);
Console.WriteLine($"{orderIds[i]} ({statuses[i]}) -> partition {result.Partition.Value}, offset {result.Offset.Value}");
}
Because the partitioner computes hash(key) % partitionCount and the key is identical for the first three sends, all three ORD-48213 events report the same partition. With the default partitioner and 6 partitions the output is (offsets depend on what is already in each log):
ORD-48213 (Created) -> partition 3, offset β¦
ORD-48213 (Paid) -> partition 3, offset β¦
ORD-48213 (Shipped) -> partition 3, offset β¦
ORD-90110 (Created) -> partition 0, offset β¦
The consumer confirms it from the other direction: whoever reads partition 3 β Consumer B in the diagram above β sees the three ORD-48213 records in production order, because within one partition Kafka preserves write order. What this demonstrates is the mechanism that guarantee depends on β keying by OrderId is what physically colocates an order's history on one ordered log instead of scattering it.
β οΈ Common Mistake: Assuming that because all three events are in the topic, a consumer reading the whole topic sees them in produced order. ORD-48213's events (all on partition 3) are always read in order by whoever owns partition 3, but that says nothing about the relative timing of records on other partitions β the topic as a whole has no global order, only partition-local order.
βοΈ Exercise: trace the worked example
Using the output and the two-consumer split above:
- Which consumer prints the
ORD-90110record, and which prints theORD-48213records? - A third instance joins the group and the assignor gives it partitions 2 and 3 (A keeps 0β1, B keeps 4β5). You produce a fourth
ORD-48213event,Delivered. Who reads it, and is it read afterShipped? - Would your answer to 2 change if the topic had been grown to 12 partitions between
ShippedandDelivered?
Check your answer
- Consumer A reads
ORD-90110(partition 0); Consumer B reads the threeORD-48213records (partition 3). - The new instance now owns partition 3, so it reads
Deliveredβ afterShipped, because both are in partition 3 in write order. (If B had not committed its progress before the rebalance, the new owner may re-read some of the earlier records first; "Consumer Groups & Offsets" covers that.) - Partly.
ORD-48213's CRC32 gives 3 for both 6 and 12 partitions, soDeliveredstill goes to partition 3 and is still read afterShipped. But adding partitions triggers a rebalance, so partition 3 may move to another instance (the range assignor now hands out 0β3, 4β7 and 8β11, so the instance that had 2β3 would get 4β7). Many other keys are not so lucky β see "What Breaks When You Add Partitions".
Inspecting Partitions and Leaders with IAdminClient
Sometimes you need to inspect a topic's structure at runtime β confirming partition count before deciding how many consumer instances to run, or checking which broker leads a partition during troubleshooting. IAdminClient.GetMetadata does this.
using System.Linq;
using Confluent.Kafka;
var adminConfig = new AdminClientConfig { BootstrapServers = "localhost:9092" };
using var admin = new AdminClientBuilder(adminConfig).Build();
// Fetch metadata for a single topic.
// (The GetMetadata(TimeSpan) overload, without a topic name, returns all topics.)
Metadata metadata = admin.GetMetadata("order-events", TimeSpan.FromSeconds(10));
var topicMetadata = metadata.Topics.Single(t => t.Topic == "order-events");
Console.WriteLine($"Topic: {topicMetadata.Topic}");
foreach (var partition in topicMetadata.Partitions)
{
Console.WriteLine(
$" Partition {partition.PartitionId}: " +
$"leader broker id = {partition.Leader}, " +
$"replicas = [{string.Join(",", partition.Replicas)}]");
}
This prints one line per partition showing the current leader (the broker producers always write to, and consumers read from by default) and the replica set, tying back to the leader/replica structure from "Anatomy of a Topic." It's a genuinely useful debugging habit: if a consumer group is reading unevenly, checking GetMetadata first tells you the actual partition count and replica layout rather than relying on documentation that may have drifted.
β οΈ Common Mistake: Using GetMetadata with a topic name as an existence check. On a cluster with auto.create.topics.enable left on, requesting metadata for a topic that doesn't exist can create it β with the broker's default partition count and replication factor, which is almost never what you wanted. When you're probing rather than confirming, use DescribeTopicsAsync(TopicCollection.OfTopicNames(...)), which reports a missing topic as an error, or list all topics with GetMetadata(TimeSpan) and look for the name.
π‘ Pro Tip: GetMetadata is synchronous and relatively cheap against the broker's cached cluster state, so it's fine in a startup health check or a diagnostic endpoint β but it's not a substitute for monitoring tooling and shouldn't be polled in a tight loop.
Taken together, these three pieces β a producer reporting where it wrote, a consumer reporting where it read, and an admin client reporting the topic's actual layout β give you a verifiable loop. You're no longer taking the partitioner on faith; you can watch a key resolve to a partition number in your own terminal.
Common Pitfalls: Partition Count, Skew, and Remapping
Every mistake in this section traces back to one habit: treating partition count and key choice as implementation details you set once and forget. Both decisions ripple through a topic's lifetime β they determine how much parallelism your consumers can ever have, what happens when traffic grows, and whether "same key, same partition" still holds after a change that felt harmless.
The Partition Count Trade-off
Partition count is the ceiling on how many consumers in a single group can process a topic in parallel, because each partition is assigned to exactly one consumer in a group at a time. That ceiling cuts both ways.
Too few partitions caps parallelism directly. If an orders topic has 3 partitions and you scale your consumer group to 6 instances hoping to double throughput, three instances sit idle β there's no partition left to assign them. You've paid for compute that does nothing.
Too many partitions looks free until you account for what each costs the brokers. Every partition is a set of log segment files the broker holds open file handles for, plus metadata the controller tracks. Each partition also needs leader election when a broker fails or restarts, and a cluster with tens of thousands of partitions makes those elections slower and more disruptive because there's more state to reconcile. There's no single correct number; it depends on target throughput per partition, consumer count, and how much broker overhead your cluster is provisioned to absorb.
Too few partitions: Right-sized for peak: Too many partitions:
[P0]-> C1 [P0]-> C1 [P0..P29] spread across
[P1]-> C2 [P1]-> C2 6 consumers, but brokers
[P2]-> C3 [P2]-> C3 now track 5x the log
C4 idle [P3]-> C4 files, leaders and
C5 idle [P4]-> C5 metadata of the right-
C6 idle [P5]-> C6 sized design (30 vs 6)
π― Key Principle: Partition count sets the maximum parallelism a consumer group can ever reach, but every partition is standing broker overhead whether or not it carries meaningful traffic. Size for expected peak consumer parallelism, not for a round number.
β οΈ Common Mistake: Picking a partition count that exactly matches today's consumer instance count, leaving no room to scale later without a repartitioning operation. Since adding partitions has consequences (below) and removing them is impossible β Kafka supports increasing partition count and offers no operation to decrease it, so shrinking means creating a new topic and migrating β pad the initial count modestly above your near-term target.
What Breaks When You Add Partitions
This is the pitfall that catches teams who understand the partitioner but haven't connected it to a live topic. The default partitioner routes with hash(key) % partitionCount, which is deterministic only while partitionCount stays fixed. The moment you add partitions β say 6 to 12 β the modulus changes for every key in your system, and each key's partition is recomputed as soon as producers see the new count (librdkafka refreshes topic metadata every 5 minutes by default, or sooner on errors).
Concretely: with the default partitioner, ORD-4471 goes to partition 3 of 6 (CRC32 % 6 = 3) but to partition 9 of 12 (CRC32 % 12 = 9). Not every key moves: going from 6 to 12, a key keeps its partition exactly when hash % 12 is below 6, which is true for about half of all keys (ORD-48213 stays on 3, for example). There is no predicting which half without computing the hash. Kafka does not rehash and move existing data: records already on partition 3 stay on partition 3. What changes is where new messages with a moved key are routed.
This matters enormously if your application depends on per-key colocation β often relied on for per-key event order or for co-locating related state in a stream-processing job. After a partition count change that guarantee breaks silently for every key that moved: old events for ORD-4471 sit on partition 3, new events appear on partition 9, and any consumer logic assuming "everything about this order is on one partition" is now wrong with no error raised anywhere.
// Adding partitions to an existing topic β this is a one-way operation.
// There is no CreatePartitions equivalent that decreases the count.
// After this call, hash(key) % partitionCount is recomputed with the new count
// (about half of all keys move when doubling), and Kafka does NOT move
// existing records to match the new mapping.
using var adminClient = new AdminClientBuilder(new AdminClientConfig
{
BootstrapServers = "localhost:9092"
}).Build();
await adminClient.CreatePartitionsAsync(new[]
{
new PartitionsSpecification
{
Topic = "order-events",
IncreaseTo = 12 // was 6 β about half of all keys now map to a different partition
}
});
β οΈ Common Mistake: Assuming increasing partition count is purely additive and low-risk. It is additive for throughput headroom, but it's a breaking change for any consumer or downstream job relying on same-key-same-partition behavior across the transition. The Java admin API's documentation warns about it (Confluent.Kafka's CreatePartitionsAsync docs do not), but the call itself succeeds, nothing is migrated, and no error appears later β the responsibility falls entirely on whoever runs the change to know which consumers depend on colocation first.
π‘ Mental Model: Partition count is closer to a database's shard count than a queue's worker pool size. Resharding a database moves data; increasing Kafka partitions does not β it changes the routing table for records not yet written, leaving old records where they are.
βοΈ Exercise: plan a partition increase (harder)
order-events must grow from 6 to 12 partitions. Using the CRC32 values from the earlier exercise:
- Which of
ORD-4471,ORD-48213,ORD-90110andORD-10233change partition? - A downstream consumer rebuilds each order's state by applying its events in partition order. Describe what can go wrong for a moved key during the change.
- Sketch a safer way to make the change.
Check your answer
- CRC32 % 12 gives 9, 3, 6 and 2. So
ORD-4471(3 β 9) andORD-90110(0 β 6) move;ORD-48213(3 β 3) andORD-10233(2 β 2) stay. Two of four β about half, as expected. - For
ORD-4471, old events stay on partition 3 and new ones go to partition 9. The two partitions are read independently, possibly by different consumers, so a newShippedevent on partition 9 can be applied before an older, still-unreadPaidevent on partition 3 β the order's history is now split across two unordered logs. - One common approach: create a new 12-partition topic; pause (or drain) the producers; let consumers finish the old topic; then switch producers and consumers to the new topic. Alternatives are dual-writing during a cut-over window, or making consumers tolerate out-of-order events (version numbers per order). In every case, decide before running
CreatePartitionsAsync, because it cannot be undone.
Skewed Keys and Hot Partitions
Even a generously sized topic bottlenecks if its keys are unevenly distributed. Key skew is when a few key values account for a disproportionate share of traffic, so the partitions those keys hash to become hot partitions β carrying far more volume than their siblings while the rest sit comparatively idle.
Imagine a separate 12-partition topic of order events keyed by CustomerId. If one enterprise customer generates a large share of total volume β very plausible in B2B β every event for that customer routes to one partition, exactly as designed. The other 11 handle everyone else comfortably, but the hot one becomes the throughput ceiling: its consumer works as fast as it can while sibling consumers finish and wait. Aggregate partition count looks sufficient on paper; achievable throughput is bounded by the busiest single partition, not the average.
This is a distinct failure from having too few partitions β adding more doesn't help if the skewed key still hashes to one of them. Mitigations are covered in "Choosing a Key: Three Questions" below. What matters here is recognizing the symptom: one partition's consumer lagging while others are caught up usually points to key skew rather than insufficient total capacity β check that partition's incoming rate to confirm, because a stuck consumer or a poison record on that partition produces the same lag picture.
π‘ Real-World Example: A dashboard showing consumer lag per partition is the fastest way to catch this. If one or two partitions climb steadily while the rest hover near zero, and those partitions also receive far more records per second than the others, that's a hot key β adding partitions or consumers won't move that line down. Topic-level average lag hides it completely.
Ordering Is Per Partition, Not Per Topic
A closely related misreading shows up constantly in teams migrating from Azure Service Bus, where people often assume a queue delivers in one global order (Service Bus only guarantees FIFO within a session). Kafka's ordering guarantee applies within a partition β records written to the same partition are read back in write order β with no guarantee about interleaving across partitions.
β Wrong thinking: "I published Order Created, then Order Shipped,
to the 'orders' topic, so consumers will always see Created before
Shipped."
β
Correct thinking: "Created and Shipped will arrive in that order
ONLY if both events land on the same partition β and keying both
by OrderId is what guarantees that. If they land on different
partitions, no ordering is guaranteed between them."
This is why key choice and ordering are inseparable: keying both events by OrderId is what forces them onto the same partition and therefore into a guaranteed sequence. The remaining caveats β producer retries ("Core Kafka Mental Model") and consumer rebalances ("Consumer Groups & Offsets") β are covered in those lessons. The instinct to fix here is narrower but non-negotiable: never reason about ordering at the topic level, only at the partition level.
Null and Accidentally-Constant Keys
The last pitfall is the quietest, because it throws no exception and survives code review β it shows up weeks later as a lopsided lag dashboard. A null key spreads records across partitions (sticky, switching every few milliseconds under the default partitioner), which is generally healthy for distribution. The danger is a key that is not null but constant: every message carries a key, nothing about the produce call looks wrong, but the key never varies.
// Bug: keying every message by a hardcoded string instead of the
// per-order identifier. Every message hashes to the SAME partition,
// no matter how many partitions the topic has.
await producer.ProduceAsync("order-events", new Message<string, string>
{
Key = "order-event", // constant literal, not OrderId
Value = JsonSerializer.Serialize(orderEvent)
});
// Fix: key by the field that actually varies per message.
await producer.ProduceAsync("order-events", new Message<string, string>
{
Key = orderEvent.OrderId, // varies per order, spreads across partitions
Value = JsonSerializer.Serialize(orderEvent)
});
In the buggy version, hash("order-event") % partitionCount evaluates to the same index for every message the producer sends β whatever the partition count, the topic degrades to the throughput of one partition; on the 6-partition order-events topic with one consumer per partition, 5 consumers sit permanently idle. This is easy to introduce when a key is built from a template string, a fixed event-type label, or a placeholder that never got replaced.
β οΈ Common Mistake: Copy-pasting a producer example that uses a literal string as the key for illustration, and shipping it without swapping in an actual per-entity field. Kafka doesn't reject a constant key β it's a perfectly valid key value β so nothing fails at compile time, deploy time, or under moderate test load. The collapse only becomes visible once real volume accumulates on the one partition it always hashes to.
Choosing a Key: Three Questions
Everything in this section comes back to the key, so here is a short, practical way to choose one. Ask three questions, in order:
- What must be ordered? Key by the entity whose events must be applied in sequence. For order events that is usually
OrderId; if a consumer needs all of a customer's orders in sequence (a credit limit, say), it isCustomerId. - How many distinct values are there, and how uneven? A key needs many more distinct values than partitions, and no single value that dominates traffic.
CustomerIdin a B2B system with one giant customer fails this;OrderIdrarely does. - Will every producer compute it the same way? Same field, same bytes (watch string casing and serialization), same partitioner on every client stack.
When questions 1 and 2 pull in opposite directions β a hot customer whose events must stay ordered β the usual escape is to order by something narrower (per order rather than per customer) and handle cross-order rules in the consumer. "Salting" a hot key (appending a suffix such as -0 to -3 to spread it over several partitions) does restore throughput, but it deliberately gives up ordering for that key, so use it only for events that don't need it.
βοΈ Exercise: pick the key
A payments topic carries PaymentAuthorized, PaymentCaptured and PaymentRefunded events. A refund must never be applied before its capture. The team proposes three keys: MerchantId (one merchant produces 40% of traffic), PaymentId, or null "for the best distribution". Which do you choose, and why not the other two?
Check your answer
PaymentId. The ordering requirement is per payment (capture before refund), and payment IDs are numerous and evenly spread.
MerchantIdwould also keep each payment's events in order, but the 40% merchant makes one partition carry 40% of all traffic β a hot partition that more partitions cannot fix.nullgives good distribution but no ordering at all: a refund and its capture can land on different partitions.
Summary and Quick Reference
You've moved from "Kafka is another message queue" to seeing it as a set of partitioned, append-only logs distributed across brokers, with keys as the routing input deciding which log a record joins. Before the deeper guarantees built on this model, lock the vocabulary down β these six terms get used loosely, and imprecise use is where a lot of production confusion starts.
The Six Terms, Precisely
A topic is the logical name producers and consumers agree on β no inherent order, no single physical location. Order and location live one level down in the partition: an ordered, immutable append-only log, the actual unit Kafka stores, parallelizes, and routes to. Within a partition, a record's position is its offset, a per-partition integer meaningless outside it β offset 500 in partition 0 and offset 500 in partition 1 are unrelated records. Physically, each partition is stored as one or more replicas (copies of its log) on different brokers (Kafka servers): one replica is the leader, the others are followers copying it. And the key is the data attached to a record that the partitioner uses to pick a partition.
| Term | What it is | What it is NOT |
|---|---|---|
| π Topic | A named, logical stream that groups related partitions | β Not itself ordered, not a single log |
| π Partition | An ordered, immutable append-only log; the real unit of storage and parallelism | β Not ordered relative to other partitions |
| π― Offset | A record's position within one specific partition | β Not a global sequence number across the topic |
| π§ Broker | A Kafka server that hosts partition replicas | β Not a unit of routing or ordering |
| π‘οΈ Replica | One copy of a partition's log β the leader or a follower β kept for fault tolerance | β Not extra parallelism: more copies, same partition count |
| π§ Key | The routing input a producer supplies per record | β Not required (except on compacted topics, where the broker keeps the latest record per key and rejects null keys) |
π― Key Principle: Ordering and parallelism belong to the partition and fault tolerance to its replicas β order lives in each partition, parallelism is capped by partition count, fault tolerance comes from replicas, and none of the three belongs to the topic name itself.
The Partitioner Mechanism, Recapped
Routing has two paths. For a keyed record, the default partitioner hashes the key and reduces it modulo the partition count β hash(key) % numPartitions β deterministic for a fixed count, so the same key always lands on the same partition. For an unkeyed record, librdkafka's default consistent_random uses a sticky random partition that changes every sticky.partitioning.linger.ms (10 ms by default); the Java client is also sticky but switches after about a batch's worth of bytes.
And the divergence worth carrying out of this lesson: librdkafka hashes keys with CRC32, the Java client with murmur2. If a topic has producers on both stacks, set Partitioner = Partitioner.Murmur2Random on the .NET side β ideally when the topic is created β or the two will disagree about where most keys belong.
using Confluent.Kafka;
var config = new ProducerConfig
{
BootstrapServers = "localhost:9092",
// Needed if a JVM client also produces to this topic. Choose it when the
// topic is new: switching an existing .NET producer over moves most keys.
Partitioner = Partitioner.Murmur2Random
};
using var producer = new ProducerBuilder<string, string>(config).Build();
// Keyed: hash(key) % partitionCount picks the partition deterministically
var keyedResult = await producer.ProduceAsync("order-events",
new Message<string, string> { Key = "ORD-4471", Value = "OrderCreated" });
Console.WriteLine($"Keyed message -> partition {keyedResult.Partition.Value}");
// Null key: goes to the current sticky partition (a random pick that changes every few ms)
var unkeyedResult = await producer.ProduceAsync("order-events",
new Message<string, string> { Value = "HeartbeatPing" });
Console.WriteLine($"Unkeyed message -> partition {unkeyedResult.Partition.Value}");
Run this twice with the same key and the keyed message reports the same partition both times (given a stable partition count), while the unkeyed message's partition isn't predictable from the call site.
π‘ Mental Model: The key is an address label and the partitioner a deterministic sorting machine β same label, same bin, every time, as long as the number of bins doesn't change and everyone feeding the machine uses the same hash function.
The Operational Facts, in One Line Each
Three things from "Common Pitfalls" are worth keeping in short-term memory while you design a topic. More partitions buy consumer parallelism but cost broker memory, file handles, and leader-election overhead, so size deliberately rather than maximize. Partition count can be increased but never decreased. And because the partitioner's modulo depends on the current count, adding partitions recomputes hash(key) % numPartitions for every key β about half of them move when you double the count β and silently breaks the colocation guarantee for those keys.
β οΈ Common Mistake: Treating partition count as something you can freely tune later the way you'd scale a stateless web service. For Kafka, changing it retroactively changes routing for a large share of existing keys, not just future capacity.
Where the Deeper Guarantees Live
Everything here has been mechanism and vocabulary. These lessons build directly on it:
- Consumer Groups & Offsets takes the "one partition, one consumer per group" rule further: how assignment and rebalancing work, and what a consumer re-reads after a crash or a move.
- Raw Confluent.Kafka Client covers the producer and consumer configuration around the partitioner β batching, delivery reports and error handling.
- Delivery Semantics & Schema Strategy goes deeper into idempotence and transactions β the producer-side conditions under which per-partition order actually holds.
π‘ Pro Tip: When designing a new topic, walk down the vocabulary table β name the topic, decide the partition count with the cost trade-off in mind, set the replication factor (3 is the common production baseline), pick a key with the three questions from "Choosing a Key" in mind, confirm the partitioner setting if anything else produces to the topic, and only then write the producer code.
Practical Next Steps
Four concrete actions. First, before writing your next producer, decide explicitly whether each message needs a key at all β a null key opts into random placement and gives up colocation, fine for independent events and wrong for anything needing per-entity ordering. Second, check whether any topic your .NET services produce to is also written by a JVM-based producer; if so, that's a Partitioner setting to reconcile today, not after an incident. Third, use DeliveryResult.Partition as a sanity check that your key is producing the colocation you expect rather than assuming it. Fourth, treat partition count as a decision made once, and if you suspect you'll need to change it, plan the migration first β work out which consumers depend on per-key colocation, as in the partition-increase exercise above.