Consumer Groups & Offsets

One partition is consumed by exactly one group member at a time. Offsets are your real checkpoints β€” not message acknowledgements.

Last generated

Lesson 3 of 10 available15 practice questions

SPACED REPETITION Β· 15 practice questions

Make this lesson stick.

Try 3 questions now. No account needed. Sample answers aren't saved.

By the end of this lesson you will be able to:

  • predict how a group's partitions are divided for any number of instances, and when extra instances sit idle;
  • explain why a committed offset is a bookmark (the next offset to read) rather than a per-message acknowledgement;
  • trace a crash through a consumer loop and say which records are lost, replayed, or processed once β€” for the .NET defaults and for the two correct commit pairings;
  • write a minimal Confluent.Kafka consumer loop with deliberate commits, narrow exception handling and a clean Close();
  • recognize five common consumer pitfalls from their symptoms.

Before you start: this lesson assumes the log, offset and partition ideas from "Core Kafka Mental Model" and "Topics, Partitions & Keys". Each major section includes a short exercise β€” try it before opening the answer.

Watch the Lesson Video28 min
The whole lesson, narrated: why consumer groups exist, one partition per consumer, offsets as checkpoints and the .NET defaults, a minimal consumer, and five pitfalls

Why Consumer Groups Exist: Scaling Reads Beyond a Single Consumer

Imagine an order-processing service reading from a Kafka topic that receives a steady stream of new orders. At launch, one consumer instance handles the load comfortably. Then a promotion goes live, order volume triples, and the same consumer starts falling further behind with every poll β€” the backlog grows, downstream fulfillment slows, and nobody touched the code. The topic didn't change shape; the read side simply couldn't keep pace. Why does adding more consumers fix this when the data lives in one place? And why does Kafka split the work by partition instead of handing out individual messages the way a traditional task queue would? Those two questions are what this section answers, by introducing the consumer group β€” the mechanism Kafka uses to let multiple consumers cooperatively read a topic, with each partition assigned to exactly one of them at a time.

The Single-Consumer Bottleneck

A Kafka topic is not one continuous stream β€” it is split into partitions, each an independently ordered, append-only log. A single consumer instance that subscribes to a topic will, by default, be handed every partition and must read them all itself. The client prefetches records on background threads, but your application's poll loop is still sequential: it takes a record, your handler runs, and only then does the loop ask for the next one. If message volume outpaces how fast that loop can process each record, the consumer's read position falls further and further behind the producer's write position. That gap is consumer lag, and it is the direct symptom of a throughput mismatch between one reader and a topic being written to faster than it's being drained.

The order-processing example makes this concrete. Suppose the orders topic has eight partitions and, during a peak sales event, producers write far more orders per second than the single consumer instance can validate, persist, and forward to a fulfillment queue. Vertical scaling β€” giving that one process more CPU or memory β€” helps only if the bottleneck is raw compute, and even then it hits a ceiling: a single thread pumping through a poll loop can only do so much sequential work per second. The structural fix is to have more than one consumer instance sharing the workload, each responsible for a subset of the partitions. That's the problem consumer groups exist to solve.

Defining the Consumer Group

A consumer group is a named set of consumer instances that coordinate to divide the partitions of one or more topics among themselves, so that the group as a whole reads the full topic while each member reads only a slice of it. Instances declare their membership by configuring the same GroupId; Kafka uses that shared identifier to know which consumers should split ownership together.

🎯 Key Principle: A consumer group parallelizes reads by dividing partitions, not by dividing individual messages. This is the single most important structural fact to internalize before anything else about groups makes sense.

That distinction matters because it's easy to picture a group the way you'd picture a thread pool pulling tasks off a shared queue β€” any idle worker grabs the next available item. Kafka consumer groups don't work that way. Ownership is assigned at the partition level: a given partition is handed to exactly one consumer instance in the group, and that instance reads every message in that partition in order, for as long as it holds that assignment. The exact rule governing assignment, and what happens when the number of consumers doesn't match the number of partitions, is the focus of the next section β€” this one only needs you to understand that partitions, not messages, are the unit that gets distributed.

πŸ’‘ Mental Model: Think of a topic's partitions as a stack of ordered logbooks, and a consumer group as a small team assigned to read the whole stack. Instead of the team fighting over which page to read next, each person is handed entire logbooks to read cover to cover. Scaling the team means handing out logbooks differently, not handing out pages.

Scaling the Order-Processing Example

Returning to the scenario: with eight partitions on the orders topic, a single consumer instance owns all eight. If that team places three more instances into the same consumer group (same GroupId, subscribed to the same topic), Kafka's group coordination redistributes the eight partitions across the four instances β€” two apiece. Each instance now keeps pace with a quarter of the traffic a single reader previously had to absorb, which is why adding instances to a group is the standard scaling lever for read throughput.

Here is a minimal illustration in .NET using the Confluent.Kafka client. Building a full poll loop is the subject of "Building a Minimal Consumer with Confluent.Kafka" later in this lesson; this snippet shows only the piece relevant here β€” the shared GroupId that makes multiple instances cooperate as one group:

using Confluent.Kafka;

// Every instance of the order-processing service uses this same GroupId.
// That shared value is what tells Kafka these instances belong to one group
// and should split partition ownership rather than each reading everything.
var config = new ConsumerConfig
{
    BootstrapServers = "localhost:9092",
    GroupId = "order-processing-service",
    AutoOffsetReset = AutoOffsetReset.Earliest // where to start with no committed offset yet (explained later)
};

using var consumer = new ConsumerBuilder<string, string>(config).Build();
consumer.Subscribe("orders");

// Each running instance of this same code, pointed at the same GroupId,
// receives a subset of the "orders" topic's partitions.

Deploy four copies of that process and you have a four-member consumer group. Kafka's group coordination (sketched below) assigns partitions across those four members automatically β€” your application code doesn't choose which partitions it gets. What matters here is the shape of the outcome: eight partitions split evenly across four readers rather than one reader draining all eight serially.

πŸ€” Did you know? A GroupId is just a string, and Kafka treats every distinct string as a separate, independent group. Two services can subscribe to the same topic under different GroupId values, and each group independently receives a full copy of every message β€” the groups share no partition assignments or offsets, so neither takes messages away from the other. That's what lets an order-processing group and a fraud-detection group both read the same orders topic in full, each at its own pace.

What This Lesson Owns, and What Comes Next

This lesson establishes the mental model you need before touching deeper mechanics: the partition-ownership rule that governs how work is divided, the distinction between an offset as a checkpoint versus a per-message acknowledgement, a minimal working .NET consumer, and the mistakes that commonly break it.

What it does not cover in depth is the machinery behind the division: the full commit API, and the triggers and consequences of a rebalance (reassigning partitions when members join or leave) β€” the "Consumer & Manual Commits" lesson picks those up.

One fact about that machinery is worth having now, because it is easy to get wrong. With the classic group protocol β€” still the default (group.protocol=classic) in librdkafka, the C library that Confluent.Kafka wraps β€” a broker acting as the group coordinator tracks membership, but the assignment itself is computed by one of the consumers, elected as group leader, using partition.assignment.strategy (librdkafka default range,roundrobin). Kafka 4.0 made the newer consumer protocol (KIP-848) generally available: with GroupProtocol = GroupProtocol.Consumer the broker computes the assignment itself and rebalances are incremental rather than pausing the whole group (the classic protocol gets incremental rebalances only from the cooperative-sticky strategy). librdkafka's documentation says its default will switch to it in a future release. Either way, the outcome this lesson relies on is the same: each partition goes to exactly one member.

Absorbing all of that alongside the basic ownership model tends to blur the distinction that makes consumer groups tractable: first understand what gets divided and why it scales reads, then learn how Kafka performs that division.

πŸ’‘ Real-World Example: A payment-notification service reading from a single-partition topic gains nothing from adding a second consumer instance to its group β€” with one partition, one instance holds it and the other sits idle. This previews the ownership rule the next section covers in full, but it's worth flagging now: scaling a consumer group only helps up to the number of partitions the topic has, so partition count is a capacity decision made well before you write any consumer code.

Why Partitions, Not Messages

It's worth pausing on why Kafka divides by partition instead of by message, since the alternative feels more intuitive if you've worked with queues where any free worker grabs the next message. A partition is an ordered log, and Kafka guarantees ordering only within a partition. If individual messages were handed to whichever consumer happened to be free, two messages written next to each other in the same partition could be processed by two instances in an unpredictable order β€” one might finish before the other started. For many use cases that's tolerable, but for anything where sequence matters within a logical unit (all events for one order ID, keyed to the same partition), losing that ordering would silently change your application's correctness. Assigning whole partitions to single consumers preserves in-partition ordering as a hard guarantee even while parallelizing across the group.

⚠️ Common Mistake: Assuming that adding consumer instances always increases throughput. A team running eight instances against an eight-partition topic doubles the fleet to sixteen under load and sees no improvement at all β€” eight of the sixteen members get zero partitions and idle. The ceiling is set by partition count, a trade-off examined in depth in the next section.

This section's job was narrow: show why a single consumer stalls against real throughput, define the consumer group as a named set of consumers splitting partition ownership, and make clear that the thing being split is partitions rather than messages. Holding onto that picture β€” a fixed number of ordered logs handed out whole to group members β€” is what makes every subsequent detail about assignment rules, rebalances, and offset commits click into place instead of feeling like arbitrary broker behavior.

✍️ Exercise: predict the effect of each change

The orders topic has 8 partitions. Group order-processing-service runs 1 instance and is falling behind. For each change, say what happens to that group's backlog.

  1. Three more instances start with the same GroupId.
  2. Instead, the fraud team starts 2 instances of a new service with GroupId = "fraud-detector" on the same topic.
  3. A separate payment-notifications topic has 1 partition, and its single consumer is also lagging. The team adds a second instance to that group.
Check your answer
  1. The 8 partitions are redistributed, 2 per instance. Each instance now drains a quarter of the traffic, so the backlog shrinks β€” assuming the bottleneck was the consumer, not a shared downstream database.
  2. Nothing changes for order-processing-service. fraud-detector is a different group: it gets its own assignment of all 8 partitions and its own offsets, and reads every message independently. It adds read load on the brokers; it relieves no one.
  3. Nothing improves. With one partition, one member owns it and the other instance sits idle. More partitions are the only way to get more parallel consumers in that group.

The Partition Ownership Rule: One Partition, One Consumer

Once a group exists, something has to decide who reads what. Kafka's answer is a single, unbending rule that shapes every capacity decision you'll make around a consumer group.

The Core Rule

🎯 Key Principle: Within a single consumer group, each partition is assigned to exactly one consumer instance at a time. Two members of the same group will never be handed the same partition simultaneously.

This is an ownership assignment, not a load-balancing suggestion. If topic orders has 6 partitions and your group order-processor has 3 running instances, the 6 partitions are handed out across the 3 consumers β€” 2 apiece with any of the built-in assignment strategies, which here decide only which two partitions each consumer gets. What never happens is two instances both pulling from partition 3 at the same time. That guarantee is what lets you reason about ordering: since a partition is an ordered log and only one consumer reads it at once, message order within that partition is preserved from the consumer's point of view. Break that guarantee and you'd lose the one thing partitions promise you.

Topic: orders (6 partitions)
Group: order-processor (3 consumer instances)

Consumer A ← partition 0
Consumer A ← partition 1
Consumer B ← partition 2
Consumer B ← partition 3
Consumer C ← partition 4
Consumer C ← partition 5

Each arrow is an exclusive ownership link for as long as group membership stays stable. Exactly how the group decides which consumer gets which partitions is the group protocol sketched in the previous section, and what happens the moment membership changes is covered in "Consumer & Manual Commits" β€” for this section, treat the assignment as a given and focus on what it implies for scaling.

Over-Provisioning: Idle Consumers

The rule cuts both ways, and the first consequence surprises teams who assume that adding instances always adds throughput. Scale order-processor from 3 instances to 8 against a 6-partition topic and two of those 8 consumers get zero partitions. They're fully connected, fully healthy, sitting in the group β€” and doing nothing.

Topic: orders (6 partitions)
Group: order-processor (8 consumer instances)

Consumer A ← partition 0
Consumer B ← partition 1
Consumer C ← partition 2
Consumer D ← partition 3
Consumer E ← partition 4
Consumer F ← partition 5
Consumer G ← (idle, no partitions)
Consumer H ← (idle, no partitions)

⚠️ Common Mistake: Assuming more pods means more throughput. An order-processing service behind a Kubernetes horizontal pod autoscaler might scale from 6 replicas to 12 under load, expecting proportional speedup. With 6 partitions, 6 of the 12 replicas idle in the group, burning compute and connection overhead while contributing nothing. The fix isn't more consumers β€” it's more partitions, a topic-level change with its own trade-offs (remapping cost, ordering scope) covered in the Topics, Partitions & Keys lesson.

Those idle consumers aren't harmless bystanders either β€” they still participate in group membership, still send heartbeats (periodic "I'm alive" messages to the group coordinator), and still factor into rebalances when they join or leave. You're paying coordination overhead for zero consumption capacity.

Under-Provisioning: One Consumer, Many Partitions

The opposite imbalance is more forgiving. Run order-processor with 2 instances against those same 6 partitions and each consumer simply owns more than one: Consumer A might get partitions 0, 1, and 2, while Consumer B gets 3, 4, and 5.

A single instance handling multiple partitions doesn't need multiple threads or any special code β€” the poll loop handles it transparently. Each call to Consume() can return a message from any partition that consumer owns, interleaved in whatever order they arrive. (Confluent.Kafka's Consume is synchronous and blocking; there is no ConsumeAsync overload, which is a common assumption for developers arriving from async-first .NET libraries.)

// A single consumer instance can own several partitions at once.
// The poll loop doesn't change β€” Consume() transparently returns
// messages from whichever owned partition has data ready.
while (!cancellationToken.IsCancellationRequested)
{
    var result = consumer.Consume(cancellationToken);

    // result.Partition tells you which of this consumer's
    // owned partitions the message came from.
    Console.WriteLine(
        $"Partition {result.Partition.Value}, " +
        $"Offset {result.Offset.Value}: {result.Message.Value}");
}

The trade-off is throughput per partition, not correctness: Consumer A is serially working through three partitions' worth of traffic instead of one. If message handling is CPU- or I/O-bound and slow, a consumer owning three partitions processes each more slowly in wall-clock terms than it would owning just one, because it's dividing its attention. That's a real capacity constraint, but it's controlled degradation β€” nothing is lost or duplicated, throughput just narrows.

Groups Are Independent: Full-Topic Fan-Out

Every example so far has stayed inside one group, where partitions are divided so each message is handled once. That's easy to confuse with a traditional competing-consumers queue, where any worker can dequeue any message and the pool competes for a shared backlog. Kafka consumer groups look similar until you add a second group.

❌ Wrong thinking: "If I add a second consumer group reading the same topic, the two groups will split the partitions between them, so each group processes half the messages."

βœ… Correct thinking: Each consumer group maintains its own independent set of partition assignments and its own committed offsets. A second group reading orders doesn't compete with the first for partitions at all β€” it gets its own full ownership assignment across its own members, and reads every message in the topic from its own offset position.

Group: order-processor (3 instances) β€” reads all 6 partitions
Consumer A ← partitions 0, 1
Consumer B ← partitions 2, 3
Consumer C ← partitions 4, 5

Group: fraud-detector (2 instances) β€” reads all 6 partitions, independently
Consumer X ← partitions 0, 1, 2
Consumer Y ← partitions 3, 4, 5

Every message published to orders is delivered once to order-processor's collective assignment and once, separately, to fraud-detector's. Neither group knows the other exists. This is what makes Kafka useful for fan-out patterns β€” feeding both an order-fulfillment pipeline and a fraud-analysis pipeline from the same stream, without duplicating the topic or the producer.

πŸ’‘ Mental Model: Go back to the stack of logbooks. Each consumer group is a separate team reading the same stack, with its own bookmarks. Within one team, every logbook is read by exactly one person at a time. A second team reads every logbook cover to cover as well, at its own pace, with its own bookmarks β€” neither team slows down or speeds up the other.

⚠️ Common Mistake: Mistaking group isolation for load-sharing across groups. It's tempting to assume adding a second group to a busy topic relieves load on the first β€” it won't, because the two groups share no offsets, assignments, or partition pool. This is different from a mistyped GroupId accidentally splitting one intended group into several unintentional full-topic readers, which is covered in "Common Pitfalls When Consuming with Confluent.Kafka."

The Capacity-Planning Trade-off

Put over- and under-provisioning together and you get the practical rule: partition count sets the hard ceiling on parallelism achievable within a single consumer group. No matter how many instances you deploy, a group can never have more actively working consumers than the topic has partitions.

This has a direct implication for topic design. If you anticipate orders eventually needing 20 parallel consumer instances at peak but provision the topic with 4 partitions, you've capped your own future scaling regardless of how much compute you're willing to throw at it. Partition count can be increased later, but it is not a transparent operation β€” it changes the mapping of keys to partitions, which breaks ordering guarantees for anything relying on key-based affinity, and it doesn't redistribute already-written data. It also cannot be decreased; shrinking means creating a new topic and migrating. That's a topic-design concern upstream of consumer behavior, but it's this ownership rule that makes the ceiling real.

πŸ”§ Scenario🎯 Partition:Consumer RatioπŸ“š Outcome
πŸ”’ Balanced6 partitions, 6 consumersEach consumer owns exactly 1 partition β€” maximum parallelism reached
πŸ”’ Under-provisioned6 partitions, 2 consumersEach consumer owns 3 partitions β€” correctness preserved, throughput narrowed
πŸ”’ Over-provisioned6 partitions, 8 consumers2 consumers sit idle β€” wasted capacity, no throughput gain

πŸ’‘ Pro Tip: When sizing a consumer group, treat partition count as the true capacity number and instance count as a variable you tune up to that ceiling β€” not past it. Provision enough partitions for the parallelism you expect at peak, then scale instances within that range as load changes, rather than repeatedly resizing the topic.

✍️ Exercise: fill in the assignment table

A topic has 12 partitions. One consumer group scales through these sizes. For each, give the largest and smallest number of partitions any instance owns, and how many instances sit idle.

Instances Most partitions on one instance Fewest Idle instances
4 ? ? ?
5 ? ? ?
14 ? ? ?

Then: why is 12 a friendlier partition count than 10 for a group that scales between 1 and 6 instances?

Check your answer
Instances Most Fewest Idle
4 3 3 0
5 3 2 0 (two instances own 3, three own 2)
14 1 0 2

With 5 instances the split cannot be even: 12 = 5 Γ— 2 + 2, so two instances carry 50% more partitions than the rest and become the slowest members. 12 splits evenly for 1, 2, 3, 4 and 6 instances β€” every size in that range except 5, as the table shows. 10 splits evenly only for 1, 2 and 5, leaving 3, 4 and 6 instances uneven. Partition counts with many divisors (6, 12, 24, 60) give you more even steps as you scale.

Offsets Are Checkpoints, Not Per-Message Acknowledgements

It's tempting to picture a Kafka consumer working like a message queue with acknowledgements: fetch a message, process it, "ack" it, move on. That model is wrong in a way with real consequences the first time your consumer crashes mid-batch. Kafka doesn't track the status of individual messages at all. It tracks a single number per partition, per consumer group, and understanding exactly what that number means β€” and when it moves β€” is the difference between predicting your system's failure behavior and being surprised by it.

Two Different Numbers Called "Offset"

The word "offset" gets used for two related but distinct things, and conflating them is where confusion starts.

The log offset is a message's permanent address inside a partition. Every record appended gets the next integer in sequence β€” 0, 1, 2, 3. This number is immutable and belongs to the message; it never changes no matter which consumer reads it or how often.

The committed offset belongs not to a message but to a consumer group. It's a single integer, stored per partition, answering one question: "where should this group resume reading if it starts fresh right now?" A committed offset of 47 means "this group has finished with everything through offset 46; hand it offset 47 next."

Partition 0 log:
  offset:   0    1    2    3    4    5    6    7
  message: [A]  [B]  [C]  [D]  [E]  [F]  [G]  [H]
                                     ^
                                     committed offset = 5
                                     (group resumes at message F)

The critical thing to notice: the committed offset is a bookmark, not a ledger of which individual messages succeeded or failed. Nothing anywhere says "message D was processed correctly, message C threw an exception." There's only the bookmark at position 5, silently implying that everything before it β€” A, B, C, D, and E β€” is considered done, whether or not that's actually true.

🎯 Key Principle: One Bookmark, Not Many Receipts

A traditional message queue often gives you per-message acknowledgement: you ack message 17 specifically, independent of whether 16 was acked. Kafka's commit model doesn't work that way. Committing an offset sets one shared checkpoint per partition. (In normal operation it moves forward; an explicit commit of a lower offset, or an operator's offset reset, moves it back β€” which is how deliberate replays work.) You cannot commit "message 16 succeeded, message 17 failed" β€” you can only say "advance the bookmark to 18," which implicitly claims both succeeded, or not advance it at all, which implicitly claims nothing past the last checkpoint succeeded.

This matters because a batch of messages rarely fails atomically. Suppose your consumer pulls ten messages, successfully writes eight to a database, then throws on the ninth. You can say "the first eight are done": commit the ninth message's offset (the eighth's offset + 1), then retry the ninth. What you cannot express is a gap β€” "one to eight and ten are done, nine is not." Your code decides, at the granularity of whatever offset value it chooses to commit, where the bookmark goes β€” and Kafka trusts that decision completely. It doesn't inspect your business logic or verify your database writes; it stores the number you gave it.

πŸ’‘ Mental Model: The committed offset is a bookmark in a paperback, not a set of highlighted sentences. You can't mark "I definitely read page 40 but skipped page 41." You place the bookmark after the last page you're confident you finished, and next time you start from there.

Why Commit Timing Is the Whole Game

Because the commit is just a number and Kafka enforces no relationship between it and actual processing outcomes, when your code advances it determines your delivery guarantees.

Commit the offset before processing and you're telling Kafka "I'm done with this" before you know that. If the consumer crashes while doing the real work β€” writing to a database, calling an API β€” that message is never processed, but on restart the consumer resumes from the committed offset, which already skipped past it. The message is silently lost from this group's perspective even though it's still sitting untouched in the log. That's at-most-once: each message is delivered zero or one times, never more, occasionally zero.

Commit the offset after processing completes and you avoid that loss but introduce the opposite risk. If the consumer crashes after finishing the work but before the commit lands, the next consumer to take that partition resumes from the older offset and reprocesses messages already handled. That's at-least-once: every message is delivered one or more times, never zero, occasionally more.

Commit BEFORE processing:
  commit(offset) β†’ process(message) β†’ [crash here] β†’ message lost

Commit AFTER processing:
  process(message) β†’ [crash here] β†’ commit(offset) never runs β†’ message reprocessed on restart

No timing choice eliminates both risks using offsets alone β€” that would require processing and committing to be one atomic operation spanning two independent systems (your business logic's side effects and Kafka's offset store), which is why exactly-once semantics require additional machinery rather than clever commit placement.

⚠️ The .NET Defaults Do Not Give You at-Least-Once

Most Kafka material tells you the defaults give you at-least-once. That is true of the Java client, whose auto-commit commits the offsets returned by the previous poll on the next poll β€” so a crash mid-batch leaves them uncommitted and they replay.

Confluent.Kafka behaves differently, because librdkafka splits the operation into two steps with two separate settings:

βš™οΈ SettingπŸ“¦ Default🎯 What it does
EnableAutoOffsetStoretrueStores the offset after each message (its offset + 1) the moment Consume() hands it to your code β€” before your handler runs
EnableAutoCommittrueCommits whatever is currently stored, on a timer
AutoCommitIntervalMs5000How often that timer fires

Read those together and the consequence is stark: with defaults, the bookmark can move past a message while your handler is still working on it, or before the handler has run at all. A crash in that window loses the message with no error and no trace. The defaults are not at-most-once either: offsets that were stored but not yet committed when the process died are simply lost, so everything handled since the last timer commit is replayed on restart. With the .NET defaults you can get both failure modes β€” loss of whatever was in flight, and duplicates of what was handled since the last timer commit β€” even though the code reads as though processing comes first.

Store vs Commit: The Two Correct Pairings

There are two configurations that give you the at-least-once behavior most business logic wants, and one combination that records nothing.

Pairing A β€” take full manual control. Turn auto-commit off and commit explicitly after the handler succeeds. Simple to reason about; costs a synchronous broker round trip per commit, so commit per batch rather than per record in throughput-sensitive services.

var config = new ConsumerConfig
{
    BootstrapServers = "localhost:9092",
    GroupId = "order-processing-group",
    EnableAutoCommit = false        // nothing moves unless we say so
};

// ...in the poll loop:
var result = consumer.Consume(cts.Token);
Process(result.Message.Value);      // work FIRST
consumer.Commit(result);            // then the bookmark moves

Pairing B β€” keep the background committer, control what it's allowed to commit. Leave auto-commit on but disable automatic offset storing, then call StoreOffset yourself after the handler succeeds. You get librdkafka's efficient batched background commits, and nothing is ever committed for a record your handler didn't finish.

var config = new ConsumerConfig
{
    BootstrapServers = "localhost:9092",
    GroupId = "order-processing-group",
    EnableAutoCommit = true,        // keep the background committer...
    EnableAutoOffsetStore = false   // ...but only let it commit what WE store
};

// ...in the poll loop:
var result = consumer.Consume(cts.Token);
Process(result.Message.Value);      // work FIRST
consumer.StoreOffset(result);       // marks this offset eligible to be committed

⚠️ Common Mistake: Combining EnableAutoCommit = false with StoreOffset. If EnableAutoOffsetStore is still at its default of true, StoreOffset throws a KafkaException on the first call β€” at least that is loud. If you also set EnableAutoOffsetStore = false, it runs without error, but stored offsets are only committed by the background committer (or an explicit parameterless Commit()), and with auto-commit off there is no background committer β€” so nothing is committed, and every restart replays from the group's last real checkpoint. Pick one pairing; never mix them.

Worked Scenario: Crash Mid-Batch

Walk through a concrete case. A consumer in group order-processing is assigned partition 0, sitting at committed offset 100, configured with Pairing A (EnableAutoCommit = false). It calls Consume() in a loop and pulls messages 100 through 104, handling each then committing.

Starting state: committed offset = 100

Step 1: process offset 100 β†’ success β†’ commit(101)   [committed = 101]
Step 2: process offset 101 β†’ success β†’ commit(102)   [committed = 102]
Step 3: process offset 102 β†’ success β†’ CRASH before commit(103) runs

The process dies right after finishing the work for offset 102 β€” say it wrote a database row β€” but before the commit for 103 reaches the broker. The committed offset is frozen at 102, even though offset 102's message was fully handled.

When the consumer (or a replacement in the same group) restarts and rejoins, it asks the broker for the last committed offset on partition 0, gets 102, and resumes fetching from there. It processes message 102 again β€” the very message it finished just before crashing β€” followed by 103 and 104, which it never reached.

Restart: fetch resumes at committed offset = 102

Step 1 (again): process offset 102 β†’ reprocessed (duplicate)
Step 2: process offset 103 β†’ first time
Step 3: process offset 104 β†’ first time

Notice what did not happen: Kafka lost no messages, and it did not need to know that offset 102 had already succeeded. It replayed everything from the last known bookmark forward. The duplicate is the direct consequence of the commit happening after processing but the crash landing in the gap between finishing the work and recording that fact. If this consumer's downstream effect β€” charging a payment, sending an email β€” isn't safe to run twice, this replay is where real damage happens, which is why idempotent processing (handling that has the same effect whether it runs once or twice) is a standard companion to at-least-once consumption rather than an optional nicety.

⚠️ Common Mistake: Assuming "the message was consumed" means "the message was committed." A message can be fetched, handed to your code, fully processed, and used to produce side effects β€” and still count as unconsumed from Kafka's perspective if the commit never happened. Consume() returning a message and the group's checkpoint advancing are two separate events, and the gap between them is exactly where crashes cause reprocessing.

Seeing the Two Offsets in Code

The distinction between a message's own offset and the group's committed offset is visible directly in the API:

// consumeResult.Offset is the LOG offset: this message's fixed position
// in the partition. It never changes and belongs to the message.
ConsumeResult<string, string> consumeResult = consumer.Consume(cancellationToken);

Console.WriteLine($"Log offset of this message: {consumeResult.Offset}");

// Process the message first...
ProcessOrder(consumeResult.Message.Value);

// ...then advance the GROUP's committed offset. This updates the bookmark,
// not a per-message flag. Everything up to and including this message is
// now implicitly "done" as far as the group is concerned.
// (Pairing A: use with EnableAutoCommit = false β€” see above.)
consumer.Commit(consumeResult);

And here's how to read back the group's current bookmark independent of any single message β€” useful for reasoning about where a restart will resume from:

// Ask the broker for this group's committed offset on a specific partition.
// This returns the checkpoint, not any information about individual messages.
var topicPartition = new TopicPartition("orders", partition: 0);
var committedOffsets = consumer.Committed(
    new[] { topicPartition },
    timeout: TimeSpan.FromSeconds(5));

var offset = committedOffsets[0].Offset;

// A group that has never committed on this partition returns Offset.Unset
// (-1001), not 0 β€” "no bookmark yet" is distinct from "bookmark at the start."
Console.WriteLine(offset == Offset.Unset
    ? "No committed offset β€” AutoOffsetReset decides where reading begins"
    : $"Group will resume at offset: {offset}");

Both snippets are simplified to isolate the offset concepts β€” real consumer code needs exception handling around Consume() and a clean shutdown path, which the next section covers.

❌ Wrong thinking: "I called Commit(), so Kafka now knows message X succeeded." βœ… Correct thinking: "I called Commit(), so the group's resume point moved forward. Kafka has no opinion on whether message X β€” or anything before it β€” actually succeeded."

Why This Distinction Outlasts Any Specific API

It's worth dwelling on why the bookmark model exists rather than per-message acknowledgement, because it shapes everything downstream. Tracking a single integer per partition is cheap: the broker needs no growing table of acknowledgement flags for every message ever produced, and a consumer rejoining after a rebalance only asks "where do I resume?" rather than reconciling a status history. That efficiency is precisely the trade-off β€” you gain a lightweight, scalable checkpoint, and you give up selective acknowledgement within a batch. Kafka's ordering-by-partition guarantee reinforces this: because messages within a partition are strictly ordered and one consumer processes them in that order, a single "resume from here" number fully describes progress. That same property makes the crash-and-resume scenario deterministic rather than ambiguous β€” the replacement consumer knows exactly which messages come after the committed offset, even though it can't know which were already handled.

✍️ Exercise: trace the same crash under three configurations

Partition 0 has committed offset 200. The consumer is handed offsets 200, 201 and 202 in turn. Handling 200 and 201 succeeds; the process crashes while the handler for 202 is still running. For each configuration, say where the restart resumes, and which of 200–202 end up processed zero, one or two times.

  1. Defaults (EnableAutoOffsetStore = true, EnableAutoCommit = true), and the 5-second auto-commit timer happened to fire just after 202 was handed to the handler.
  2. Defaults again, but the timer last fired just after 200 was handed out, and not since.
  3. Pairing A (EnableAutoCommit = false, Commit(result) after each successful handler).
Check your answer
  1. When 202 was handed out, its next offset (203) was stored at once, and the timer committed 203. The restart resumes at 203: 200 and 201 were processed once, 202 is lost β€” its handler never finished and nothing will retry it.
  2. The timer committed 201 (stored when 200 was handed out) and never ran again. The restart resumes at 201: 200 once, 201 twice (duplicate), 202 once (its first completed run).
  3. Commits after 200 and 201 leave the committed offset at 202. The restart resumes at 202: 200 and 201 once, 202 once β€” no loss, and no duplicate because 202 never completed before the crash. (Had the crash come after 202's handler but before its commit, 202 would run twice β€” the at-least-once duplicate.)

Cases 1 and 2 are the same code and the same crash; only the timer's timing differs. That is why the defaults are not a delivery guarantee.

Building a Minimal Consumer with Confluent.Kafka

Everything discussed so far β€” partition ownership, offset semantics β€” is abstract until you write the code that joins a group and pulls messages off a partition. This section builds that code piece by piece using the official Confluent.Kafka client. The goal isn't a production-ready service; it's a correct, minimal skeleton you'll extend in later lessons.

Configuring the Consumer

Every consumer starts with a ConsumerConfig, a strongly typed wrapper around the settings the underlying client library accepts. Five of them matter before you write a line of consuming logic.

BootstrapServers is the initial contact point β€” one or more broker addresses the client uses to discover the rest of the cluster. It needn't list every broker; the client learns full cluster metadata after the first successful connection. A phone number to dial, not the whole address book.

GroupId determines which consumer group this instance joins. This single field turns an isolated client into a member of a coordinated group β€” every consumer sharing the same GroupId against the same topic has partitions divided among them per the one-partition-one-consumer rule. Get it wrong β€” a typo, a copy-paste error, an environment suffix left in by accident β€” and you silently create a second group that reads the entire topic independently; that failure mode gets its own treatment in "Common Pitfalls."

AutoOffsetReset controls where this group starts reading a partition when it has no committed offset to resume from β€” for example, the very first time it reads that partition. Earliest starts at the beginning of the retained log; Latest starts from whatever is produced from this moment on, ignoring history. It applies only when there is no valid offset to read from β€” none was ever committed, the group's stored offsets expired, or the offset is out of range because retention has deleted those records (this can also happen to a running consumer that falls behind retention, which then jumps to Latest by default, with no exception β€” librdkafka only logs a warning); it does not override a valid checkpoint. A new group pointed at a long-running topic with AutoOffsetReset left at its default of Latest silently skips every message already in the log β€” a common surprise for anyone expecting to "replay from the start."

EnableAutoCommit and its companion EnableAutoOffsetStore decide when the bookmark moves. Both default to true, which as the previous section showed guarantees neither at-least-once nor at-most-once. Pick one of the two pairings deliberately rather than inheriting the defaults:

using Confluent.Kafka;

var config = new ConsumerConfig
{
    BootstrapServers = "localhost:9092",
    GroupId = "order-processing-group",
    AutoOffsetReset = AutoOffsetReset.Earliest,

    // Pairing A: we commit explicitly, after the handler succeeds.
    // (Pairing B is EnableAutoCommit = true + EnableAutoOffsetStore = false.)
    EnableAutoCommit = false
};

πŸ’‘ Pro Tip: Treat GroupId as a deployment-time configuration value, not a hardcoded literal. Two services with slightly different GroupId strings pointed at one topic each receive their own full copy of the data β€” sometimes exactly what you want, sometimes a bug.

Constructing the Client and Subscribing

With configuration in hand, the client is built through ConsumerBuilder<TKey, TValue>, a generic builder that also lets you specify how keys and values are deserialized. For plain string payloads the built-in Deserializers.Utf8 is enough; typed payloads (Avro, JSON, Protobuf) plug into the same builder but are outside this section's scope.

using var consumer = new ConsumerBuilder<Ignore, string>(config)
    .Build();

consumer.Subscribe("order-events");

The key type here is Ignore, meaning the code doesn't care about message keys β€” only the value matters. Calling Subscribe does not connect to a partition immediately; it registers interest in the topic and hands partition assignment to the group protocol described at the start of this lesson. For now it's enough that Subscribe is a declaration of intent, and the actual assignment arrives asynchronously once you start calling Consume.

⚠️ Common Mistake: Calling Subscribe repeatedly, expecting topics to accumulate. Subscribe replaces the current subscription set entirely rather than adding to it β€” for multiple topics, pass them in one call: consumer.Subscribe(new[] { "order-events", "order-events-retry" }).

The Poll Loop

Kafka consumers are pull-based: nothing is pushed to your process, so you must actively ask for the next message. This is the poll loop β€” a while loop repeatedly calling consumer.Consume(cancellationToken), blocking until a message is available or the token is cancelled, and returning a ConsumeResult<TKey, TValue> containing the message, its topic-partition, and its offset.

using var cts = new CancellationTokenSource();
Console.CancelKeyPress += (_, e) =>
{
    e.Cancel = true; // prevent immediate process kill
    cts.Cancel();
};

try
{
    while (!cts.IsCancellationRequested)
    {
        try
        {
            var result = consumer.Consume(cts.Token);

            Console.WriteLine(
                $"Partition {result.Partition.Value}, " +
                $"Offset {result.Offset.Value}: {result.Message.Value}");

            // Business logic goes here. Keep it fast β€” a slow handler
            // delays the next Consume() call and can trigger a rebalance,
            // covered in "Common Pitfalls When Consuming with Confluent.Kafka."

            consumer.Commit(result);   // Pairing A: bookmark moves after the work
        }
        catch (ConsumeException ex)
        {
            Console.WriteLine($"Consume error: {ex.Error.Reason}");
            // Keep the loop alive β€” but Pitfall 5 shows why a real consumer
            // must also dead-letter the record (copy it to a separate topic
            // or store for later inspection), not just log it.
        }
        catch (KafkaException ex)
        {
            // Commit failed β€” usually because a rebalance already moved the
            // partition away. The new owner resumes from the last good commit.
            Console.WriteLine($"Commit failed: {ex.Error.Reason}");
        }
    }
}
catch (OperationCanceledException)
{
    // Expected: cancelling the token unblocks Consume() by throwing.
}
finally
{
    consumer.Close();
}

Walk through what this loop does. Each Consume returns one ConsumeResult β€” there is no batch API surfaced here, even though internally the client fetches from the broker in batches. The Offset.Value printed alongside each message is the log offset, distinct from the committed offset the group records. And note the two nested try blocks: the inner one catches per-message failures so the loop survives them, while the outer one catches the cancellation that ends the loop. OperationCanceledException is not a ConsumeException, so without that outer catch it escapes the method entirely after finally runs β€” an unhandled exception on every clean shutdown. The inner block has two catches. ConsumeException comes from Consume(). A failed Commit throws a KafkaException (sometimes its subclass TopicPartitionOffsetException), most often because a rebalance has just moved the partition to another member β€” without that second catch, one routine rebalance would crash the process. (ConsumeException derives from KafkaException, so it must be caught first.)

Poll loop lifecycle, one iteration:

consumer.Consume(token) blocks
    ↓
broker returns next message (or token cancelled β†’ OperationCanceledException)
    ↓
ConsumeResult returned: { Partition, Offset, Message }
    ↓
handler code runs on the message
    ↓
offset committed (or stored, depending on pairing)
    ↓
loop repeats

Handling ConsumeException Without Crashing

ConsumeException is thrown by Consume() when a record cannot be deserialized, or when the client reports an error for a particular topic or partition (for example, the topic does not exist or you are not authorized to read it), or when the consumer itself hits a group-level error such as Local_MaxPollExceeded (a handler ran too long; see Pitfall 2). Broker connection problems are different: librdkafka retries them itself and reports them to the error handler (SetErrorHandler on the builder), not as exceptions from Consume(). If it escapes the loop, the process terminates and every partition this consumer owned must be picked up by another group member β€” or sits unprocessed until this instance restarts β€” a disruption for what might be a single corrupt record. Catching it at the point of the Consume call keeps the loop alive so the next message is still attempted.

⚠️ Common Mistake: This inner catch is deliberately narrow. It's easy to widen it into a blanket catch (Exception) that also swallows business-logic failures inside your handler, silently discarding messages that needed a dead-letter queue, an alert, or a retry. A blanket catch would also swallow OperationCanceledException. In this loop the while condition still ends it, but in a while (true) loop β€” or one that checks a different token β€” it would turn your shutdown signal into a tight spin loop. Pitfall 5 in the next section covers what even a narrow catch still has to do. The pattern here is scoped to the exceptions the client raises.

Shutting Down Cleanly

The last piece is graceful shutdown, and it matters more than it looks. A CancellationTokenSource is wired to the process's cancel signal (Console.CancelKeyPress in a console app, or a hosted service's stopping token), and cancelling it causes the blocking Consume call to throw OperationCanceledException, which the outer catch absorbs. The critical step is in finally: calling consumer.Close().

Close() does two things that matter to the group. If auto-commit is enabled (Pairing B), it commits the stored offsets one last time β€” with Pairing A there is nothing pending, because you committed after each message. And β€” more importantly here β€” it sends an explicit "leaving group" notification to the coordinator, so this consumer's partitions are reassigned to the remaining members right away. (With static membership β€” GroupInstanceId set β€” no leave is sent: the group deliberately waits for that member to come back within the session timeout, the time the coordinator waits without hearing from a member before declaring it gone.) Skip Close() and the coordinator can't know the member is gone until its session timeout elapses, leaving those partitions unread for that entire window.

🎯 Key Principle: A poll loop is only "minimal" in business logic β€” its structural requirements (deliberate commit configuration, a cancellable loop, narrow exception handling around Consume, an outer catch for cancellation, and a guaranteed Close()) are not optional extras to layer on later. Every consumer you build from here sits on exactly this skeleton.

Putting the Pieces Together

using Confluent.Kafka;

var config = new ConsumerConfig
{
    BootstrapServers = "localhost:9092",
    GroupId = "order-processing-group",
    AutoOffsetReset = AutoOffsetReset.Earliest,
    EnableAutoCommit = false          // Pairing A: explicit commits
};

using var consumer = new ConsumerBuilder<Ignore, string>(config).Build();
consumer.Subscribe("order-events");

using var cts = new CancellationTokenSource();
Console.CancelKeyPress += (_, e) => { e.Cancel = true; cts.Cancel(); };

try
{
    while (!cts.IsCancellationRequested)
    {
        try
        {
            var result = consumer.Consume(cts.Token);
            Console.WriteLine($"Offset {result.Offset.Value}: {result.Message.Value}");
            consumer.Commit(result);
        }
        catch (ConsumeException ex)
        {
            Console.WriteLine($"Consume error: {ex.Error.Reason}");
        }
        catch (KafkaException ex)
        {
            Console.WriteLine($"Commit failed: {ex.Error.Reason}");
        }
    }
}
catch (OperationCanceledException)
{
    // graceful shutdown
}
finally
{
    consumer.Close();
}

This is deliberately incomplete as a production artifact β€” no dead-letter routing, no batched commits, no protection against a slow handler starving the loop. What it gives you is a correct scaffold: every instance built this way registers with the group named by GroupId, receives a subset of the topic's partitions, records progress only after the work is done, and leaves the group cleanly on shutdown.

✍️ Exercise: review this consumer

A teammate submits this worker. Find three problems and say what each one does at runtime.

var config = new ConsumerConfig
{
    BootstrapServers = "localhost:9092",
    GroupId = "invoice-writer",
    EnableAutoCommit = false,
    EnableAutoOffsetStore = false
};

using var consumer = new ConsumerBuilder<string, string>(config).Build();
consumer.Subscribe("invoices");

while (true)
{
    try
    {
        var result = consumer.Consume(stoppingToken);
        await WriteInvoiceAsync(result.Message.Value);
        consumer.StoreOffset(result);
    }
    catch (Exception ex)
    {
        logger.LogError(ex, "Failed");
    }
}
Check your answer
  1. Mixed pairing. Auto-commit is off, so nothing ever commits the offsets StoreOffset records. Every restart goes back to the group's last real commit and rewrites every invoice since then; if the group has never committed, AutoOffsetReset applies instead, and at its default Latest each restart silently skips everything produced while the worker was down. Fix: Commit(result) (Pairing A), or EnableAutoCommit = true with StoreOffset (Pairing B).
  2. Blanket catch in while (true). When stoppingToken is cancelled, Consume throws OperationCanceledException, the catch logs it, and the loop calls Consume again, which throws again at once β€” a tight loop that never shuts down. It also hides failed invoice writes: the loop just moves on to the next record, with nothing retrying the failed one or recording it β€” and once the commit pairing is fixed, the next successful record's offset moves the bookmark past it for good. Catch ConsumeException narrowly, let cancellation end the loop, and send failed records to a dead-letter path.
  3. No Close(). On shutdown the group only notices this member is gone after the session timeout (45 seconds by default in librdkafka), so its partitions sit unread for that long on every deploy. Put consumer.Close() in a finally.

Common Pitfalls When Consuming with Confluent.Kafka

The minimal consumer works fine in a demo. Production traffic exposes a different set of failure modes. Most damage in real deployments comes from a handful of implementation habits that look harmless in isolation but break ordering, stall processing, or quietly drop data. Here are five, and how to recognize them before they reach production.

Pitfall 1: Sharing One IConsumer Across Threads

It's tempting to treat IConsumer<TKey,TValue> like any other injectable service β€” register it once, hand it to multiple worker threads, let them all call Consume(). It compiles, runs, and may even look faster β€” librdkafka's API is thread-safe, so nothing is corrupted. What breaks is your guarantees:

  • Each Consume() call hands out the next record from whichever partition has data, so consecutive records of one partition land on different threads and are processed concurrently, in no defined order.
  • With the default auto offset store, each record's offset is stored the moment it is handed out, so the background commit can move the bookmark past records other threads have not finished.
  • Rebalance and error handlers run as a side effect of Consume() on whichever thread happens to be calling it β€” while the other threads are still working on records from partitions that may be in the middle of being revoked.
// ❌ Wrong thinking: "more threads calling Consume() means faster throughput"
var consumer = new ConsumerBuilder<string, string>(config).Build();
consumer.Subscribe("orders");

// Spinning up several threads against the SAME consumer instance
for (int i = 0; i < 4; i++)
{
    Task.Run(() =>
    {
        while (!cts.IsCancellationRequested)
        {
            var result = consumer.Consume(cts.Token); // one partition's records now spread over 4 threads
            Process(result);
        }
    });
}

None of this throws β€” that's what makes it dangerous. You see events for the same key applied out of order, records that were never processed after a crash, and extra duplicates after rebalances. Because it is timing-dependent, it often passes local testing and only surfaces under load.

βœ… Correct thinking: one IConsumer instance belongs to exactly one thread's poll loop. To parallelize across cores, either run multiple consumer instances (each with its own IConsumer, all sharing a GroupId so the group splits partitions between them), or keep a single-threaded poll loop that hands each ConsumeResult to a bounded worker pool while the loop keeps polling.

using System.Threading.Channels;

// One consumer, one poll-loop thread β€” parallelism happens downstream
var channel = Channel.CreateBounded<ConsumeResult<string, string>>(capacity: 100);

// Single thread owns the consumer and only calls Consume() here
// (Close() and the minimal consumer's try/catch blocks are left out to keep this short)
_ = Task.Run(async () =>
{
    using var consumer = new ConsumerBuilder<string, string>(config).Build();
    consumer.Subscribe("orders");
    while (!cts.IsCancellationRequested)
    {
        var result = consumer.Consume(cts.Token);
        await channel.Writer.WriteAsync(result, cts.Token);
    }
});

// Separate worker tasks pull from the channel and do the actual processing
for (int i = 0; i < 4; i++)
{
    _ = Task.Run(async () =>
    {
        await foreach (var result in channel.Reader.ReadAllAsync(cts.Token))
        {
            Process(result); // the consumer is untouched here β€” but see the caveat below
        }
    });
}

The rule to internalize: one IConsumer, one poll loop β€” not because the object would break, but because ordering, offset storage and rebalance callbacks all assume one sequential reader. The messages it hands you are plain data, safe to pass to other threads once they're off the poll loop β€” subject to the caveat below.

⚠️ This hand-off pattern moves the offset problem rather than solving it. As written, it is deliberately incomplete on two fronts. First, once messages leave the poll loop, the loop keeps polling β€” so whatever advances the bookmark (auto-store, auto-commit, or a commit call in the loop) runs ahead of the workers, and a crash loses everything in flight. Second, several workers can process messages from the same partition concurrently, which discards the per-partition ordering the whole ownership model exists to protect. Making it correct means keeping each partition's records in order (for example, one worker per partition, or per key), tracking completion per partition, and storing only the offset after the highest contiguously completed message β€” real work that no single API call does for you. Use the simple version only where duplicate and out-of-order processing are both genuinely acceptable.

Pitfall 2: Long-Running Handlers Blocking the Poll Loop

Even with a correctly single-threaded consumer, it's easy to do too much work per iteration:

while (!cts.IsCancellationRequested)
{
    var result = consumer.Consume(cts.Token);
    await CallDownstreamApiAsync(result.Message.Value); // could take seconds
    await SaveToDatabase(result.Message.Value);
}

It's worth knowing exactly which liveness check this trips. Group heartbeats are sent on a background thread β€” in librdkafka, and in the Java client since Kafka 0.10.1 β€” so a blocked handler does not fail the session timeout. What it trips is MaxPollIntervalMs (default 300000 ms β€” five minutes): the maximum gap allowed between your calls into Consume(). Exceed it and the client logs Application maximum poll interval exceeded, leaves the group, and your next commit fails because you no longer own the partition; the next Consume() throws a ConsumeException with ErrorCode.Local_MaxPollExceeded, and the consumer rejoins as polling resumes. The record gets reprocessed by whoever inherits it.

The practical takeaway: keep per-iteration work short, and if a message requires slow processing, hand it to a background worker (with the caveats above) rather than doing the slow work inline. Raising MaxPollIntervalMs is a legitimate second option when work genuinely takes minutes β€” the cost is that a truly stuck consumer goes undetected that much longer.

⚠️ Common Mistake: Teams discover this only after volume grows β€” a handler taking milliseconds against test data takes seconds against a real downstream service, and a smoothly running group starts rebalancing under load, exactly when it can least afford the disruption.

Pitfall 3: Skipping Close() on Shutdown

When a consumer process exits, the group finds out in one of two very different ways. If the consumer calls Close() first, it sends an explicit "I'm leaving" notification to the coordinator, which rebalances immediately and reassigns those partitions with minimal delay. If the process exits abruptly β€” crashes, gets killed, or simply never calls Close() β€” the broker can't know until it stops seeing heartbeats, then waits out the session timeout before declaring the member dead.

// ❌ No Close() β€” broker only notices via session timeout
var consumer = new ConsumerBuilder<string, string>(config).Build();
consumer.Subscribe("orders");
try
{
    while (!cts.IsCancellationRequested)
    {
        var result = consumer.Consume(cts.Token);
        Process(result);
    }
}
catch (OperationCanceledException)
{
    // process just exits here β€” partitions sit unclaimed until timeout
}
// βœ… Explicit Close() in a finally block
var consumer = new ConsumerBuilder<string, string>(config).Build();
consumer.Subscribe("orders");
try
{
    while (!cts.IsCancellationRequested)
    {
        var result = consumer.Consume(cts.Token);
        Process(result);
    }
}
catch (OperationCanceledException)
{
    // expected on graceful shutdown
}
finally
{
    consumer.Close(); // leaves the group cleanly, triggers immediate rebalance
}

The practical cost of skipping Close() is a window during which nobody reads that consumer's partitions β€” messages queue up unprocessed until the session timeout expires. In a rolling deployment restarting instances one at a time, this repeats on every restart and adds up to real, avoidable latency. Dispose() releases the underlying client handle but does not commit pending offsets or send the leave-group notification, so a controlled shutdown wants Close() called explicitly before disposal.

Pitfall 4: Mistyped or Copy-Pasted GroupId Values

The GroupId string is the only thing telling the broker which consumers should share a topic's partitions. Because it's just a string in configuration, it's exactly the kind of value that gets copy-pasted between services, typo'd in an environment variable, or left as a tutorial placeholder.

// Service A
var configA = new ConsumerConfig
{
    BootstrapServers = "broker:9092",
    GroupId = "order-processor",
    AutoOffsetReset = AutoOffsetReset.Earliest
};

// Service B β€” meant to join the same group, but has a trailing typo
var configB = new ConsumerConfig
{
    BootstrapServers = "broker:9092",
    GroupId = "order-processor ", // trailing space β€” a DIFFERENT group to the broker
    AutoOffsetReset = AutoOffsetReset.Earliest
};

The broker treats "order-processor" and "order-processor " as two entirely distinct groups. Instead of splitting partitions between the two instances, each forms its own single-member group and reads every partition independently. No exception is thrown; both consumers appear to work and per-instance throughput looks fine β€” but downstream, every message is processed twice, and if the workload was split for capacity reasons, neither instance is relieved of any load.

πŸ’‘ Pro Tip: Because this failure produces no errors, the reliable way to catch it is to check group membership directly rather than trust the config file. kafka-consumer-groups.sh --bootstrap-server broker:9092 --describe --group order-processor --members lists the members the broker actually has in that group and how many partitions each owns (from .NET, IAdminClient.DescribeConsumerGroupsAsync returns the same information) β€” the fastest way to confirm your instances are sharing partitions as intended.

Pitfall 5: Logging and Discarding ConsumeException

The minimal consumer wraps Consume() in try/catch (ConsumeException) so one malformed message doesn't crash the process. That's a reasonable baseline, but it's easy to take the defensive instinct too far β€” logging the exception and moving on with no record of which message caused it and no path for dealing with it later.

// ❌ Wrong thinking: "as long as we don't crash, we're fine"
while (!cts.IsCancellationRequested)
{
    try
    {
        var result = consumer.Consume(cts.Token);
        Process(result);
    }
    catch (ConsumeException e)
    {
        Console.WriteLine($"Error: {e.Error.Reason}"); // and then... nothing
    }
}

The problem isn't that the process survives β€” it's that the poison message (malformed, corrupted, or otherwise undeserializable) simply vanishes from view. By the time Consume() throws, the client has already moved past that record: the next call returns the next record, and as soon as a later offset is committed or stored, the bad one is behind the bookmark for good. There is no automatic retry and no record of what went wrong. Anyone debugging a missing order later has no trail to follow.

// βœ… Correct thinking: capture context and route to a dead-letter path
while (!cts.IsCancellationRequested)
{
    try
    {
        var result = consumer.Consume(cts.Token);
        Process(result);
    }
    catch (ConsumeException e) when (e.Error.Code is ErrorCode.Local_KeyDeserialization
                                              or ErrorCode.Local_ValueDeserialization)
    {
        // A permanently bad record: preserve enough detail to diagnose and reprocess later
        await deadLetterSink.SendAsync(new
        {
            e.Error.Reason,
            Topic = e.ConsumerRecord?.Topic,
            Partition = e.ConsumerRecord?.Partition.Value,
            Offset = e.ConsumerRecord?.Offset.Value,
            // ConsumerRecord carries the raw, undeserialized key and value bytes
            Key = e.ConsumerRecord?.Message?.Key,
            Value = e.ConsumerRecord?.Message?.Value,
            Timestamp = DateTimeOffset.UtcNow
        }, cts.Token);
    }
    catch (ConsumeException e) when (e.Error.Code == ErrorCode.Local_MaxPollExceeded)
    {
        // About the consumer, not a record: log it and keep polling (see below)
        Console.WriteLine($"Consumer fell out of the group: {e.Error.Reason}");
    }
    // Any other ConsumeException escapes this loop: triage it by code (see below)
}

(Simplified β€” a production dead-letter path typically needs its own topic or durable store, and should carry the record's headers too. The habit that matters is capture and route, not log and discard.)

⚠️ Common Mistake: Treating every ConsumeException the same way. Check e.Error.Code. A deserialization failure (ErrorCode.Local_ValueDeserialization or Local_KeyDeserialization) means one permanently bad record that belongs in the dead-letter path. An authorization or unknown-topic error means the consumer itself is misconfigured and no record will succeed β€” dead-lettering and continuing would just churn through errors, so alert and stop instead. Local_MaxPollExceeded is about the consumer, not a record: log it, keep polling (the consumer rejoins the group), and fix the slow handler (Pitfall 2). (Transient broker connectivity problems never arrive here: librdkafka retries them and reports them to the error handler.)

Why These Pitfalls Compound

None of these five are exotic β€” each is a small, easy decision in otherwise ordinary consumer code. What makes them worth studying together is that most of them fail silently rather than loudly: a shared IConsumer doesn't throw "you violated thread safety"; a mistyped GroupId doesn't refuse to start; a swallowed ConsumeException doesn't page anyone. Each quietly erodes something you were relying on β€” per-partition ordering, timely reassignment, a record of every message that failed β€” without announcing it. Diagnose them by checking what actually happens β€” group membership and partition assignment, per-key ordering, the logs, the dead-letter trail β€” against what you intend, rather than trusting that the configuration you wrote matches the behavior you're getting.

Symptom you observe          Likely pitfall to check first
─────────────────────────────────────────────────────────
Every message processed      GroupId mismatch (Pitfall 4)
  twice
Unexplained rebalances       Slow handler in poll loop (Pitfall 2)
Same-key events out of order Shared IConsumer across threads (Pitfall 1)
Delayed reassignment on      Missing Close() (Pitfall 3)
  deploy or restart
Messages disappearing        Log-and-discard ConsumeException (Pitfall 5)
  with no trace
Messages lost or duplicated  Default auto-store/auto-commit pairing
  on restart, no error         (see "Offsets Are Checkpoints")

Keeping this mapping in mind turns a vague "the consumer is acting weird" investigation into a targeted check of one specific piece of code, which is usually where the fix ends up living.

✍️ Exercise: diagnose from the symptom

Name the pitfall behind each report and the first thing you would check.

  1. After every rolling deploy, lag on a few partitions climbs for about 45 seconds, then drains. Nothing else is wrong.
  2. Under a load test, the group rebalances every few minutes, and the logs contain Application maximum poll interval exceeded.
  3. Finance finds every invoice since the last deploy booked twice. Both consumer instances are healthy, and each shows every partition of the topic assigned to it.
  4. A customer's OrderCancelled is sometimes applied before its OrderPlaced, even though both carry the same key.
Check your answer
  1. Pitfall 3, missing Close(). 45 seconds is librdkafka's default session.timeout.ms: the group waits that long to notice each restarted instance has gone. Check that shutdown reaches Close().
  2. Pitfall 2, slow handler. The loop is not returning to Consume() within max.poll.interval.ms, so the client leaves the group. Measure handler time per record; move slow work off the loop or bound it.
  3. Pitfall 4, GroupId mismatch. Two instances each owning every partition means two groups. Compare the GroupId values byte for byte, and confirm with kafka-consumer-groups.sh --describe --members.
  4. Pitfall 1, or the worker hand-off without per-partition ordering: same key means same partition, so the order was lost after Consume(). Check for several threads calling Consume(), or a worker pool that doesn't keep each partition's (or key's) records on one worker.

Key Takeaways: Consumer Groups and Offsets

You've walked through why consumer groups exist, the ownership rule governing how they divide work, the difference between a message offset and a committed offset, and a working Confluent.Kafka consumer loop. Before moving into the mechanics that explain how Kafka enforces all this, it's worth consolidating the model into something you can check your own code against. This is a review pass, not new material β€” treat it as a checklist.

Recap: Ownership Caps Parallelism

The rule shaping everything else: within one consumer group, a partition is assigned to exactly one consumer instance at a time. That's what separates consumer groups from a competing-consumers queue β€” the group divides partitions, not messages. The guarantee is about assignment: a member evicted mid-handler (Pitfall 2) can still be finishing a record while the new owner starts on that partition β€” one more reason handlers should be idempotent.

The consequence is a hard ceiling: a topic with 6 partitions supports at most 6 usefully active consumers in one group. Add a 7th and one of the seven sits idle until a member leaves. Conversely, 2 consumers against 6 partitions means each owns 3 and interleaves work across them. Neither is an error β€” both are mechanical consequences of the ownership rule, and it's why partition count is the number you plan around before instance count.

🎯 Key Principle: Consumer group scaling is bounded by partition count, not by however many instances you're willing to deploy. Adding consumers past that ceiling adds no throughput β€” only idle processes waiting for a rebalance.

Recap: Offsets Checkpoint Progress, They Don't Acknowledge Messages

The second load-bearing idea is the distinction between a message's log offset β€” its fixed position in the partition β€” and the group's committed offset β€” the saved bookmark for where to resume. A commit doesn't mark individual messages done; it moves one number forward, and Kafka treats everything before it as handled, regardless of whether every message in between actually succeeded.

That's why commit timing determines your delivery semantics:

  • Commit before processing β†’ a crash mid-handler skips those messages on restart. Risks message loss.
  • Commit after processing β†’ a crash after the work but before the commit replays them. Risks duplicate processing, the foundation of at-least-once.

And the .NET-specific trap worth carrying out of this lesson: Confluent.Kafka's defaults do not give you the second shape. EnableAutoOffsetStore defaults to true, storing the next offset the instant Consume() returns a message β€” before your handler runs β€” and EnableAutoCommit commits whatever is stored on a 5-second timer. Getting at-least-once requires deliberately choosing one of two pairings:

PairingConfigProgress recorded by
A β€” manualEnableAutoCommit = falseCommit(result) after the handler
B β€” storedEnableAutoCommit = true
EnableAutoOffsetStore = false
StoreOffset(result) after the handler
❌ MixedEnableAutoCommit = false + StoreOffsetNothing β€” StoreOffset throws (auto-store still on), or stores offsets that nothing commits (auto-store off)

Neither pairing gives exactly-once for free β€” that requires idempotent handling or transactional writes, outside what committed offsets alone can guarantee. The full commit API surface β€” batching commits, committing from rebalance handlers β€” belongs to the Consumer & Manual Commits lesson; what matters here is recognizing commit timing as a design decision rather than an implementation detail you can inherit.

Recap: The Minimal .NET Consumer Shape

Stripped to essentials, a working consumer needs five pieces: a GroupId, a deliberate commit configuration, a Subscribe() call, a Consume() loop that runs until told to stop, with narrow catches so one bad record or failed commit doesn't end it, and a Close() on shutdown so the broker learns about departure immediately instead of waiting out a session timeout.

var config = new ConsumerConfig
{
    BootstrapServers = "localhost:9092",
    GroupId = "order-processing-group", // ties this instance to a specific group
    AutoOffsetReset = AutoOffsetReset.Earliest,
    EnableAutoCommit = false            // deliberate: we commit after the work
};

using var consumer = new ConsumerBuilder<string, string>(config).Build();
consumer.Subscribe("orders");

using var cts = new CancellationTokenSource();
Console.CancelKeyPress += (_, e) =>
{
    e.Cancel = true; // let the finally block run Close() instead of killing the process
    cts.Cancel();
};

try
{
    while (!cts.IsCancellationRequested)
    {
        try
        {
            var result = consumer.Consume(cts.Token);
            // handler logic goes here
            consumer.Commit(result);
        }
        catch (ConsumeException ex)
        {
            // triage by ex.Error.Code; dead-letter bad records (Pitfall 5)
            Console.WriteLine($"Consume error: {ex.Error.Reason}");
        }
        catch (KafkaException ex)
        {
            // usually a commit that failed after a rebalance; keep polling
            Console.WriteLine($"Commit failed: {ex.Error.Reason}");
        }
    }
}
catch (OperationCanceledException)
{
    // expected on shutdown β€” Consume() throws when the token is cancelled
}
finally
{
    consumer.Close(); // leaves the group cleanly rather than waiting on a timeout
}

Each piece maps back to a concept: GroupId determines which ownership set this instance joins, the commit configuration determines whether you get loss or duplicates on a crash, Subscribe opts into partition assignment, the loop is where the work happens, and Close() lets the group react immediately rather than waiting on a failure-detection window.

πŸ’‘ Mental Model: The pieces answer five questions β€” which group am I in (GroupId), when does my progress count (commit config), what am I reading (Subscribe), how do I get messages without one failure ending the loop (the loop and its narrow catches), and how do I leave politely (Close). A consumer missing any one is incomplete, not just suboptimal.

Review Checklist

Before treating a consumer as production-ready, walk six checks that map to failure modes covered above:

  1. Commit configuration β€” have you deliberately chosen Pairing A or B, or are you inheriting defaults that can both lose and duplicate messages? This is the check most likely to be wrong in code that otherwise looks correct.
  2. One poll loop per consumer β€” is a single IConsumer<TKey,TValue> instance ever consumed from more than one thread? Nothing throws, but per-partition ordering and offset correctness are gone.
  3. Handler duration β€” does the code inside the loop do slow, blocking work before looping back to Consume()? Exceeding MaxPollIntervalMs between calls drops you out of the group.
  4. Shutdown handling β€” does every exit path, including exceptions, reach Close()? (A process killed outright β€” SIGKILL, out of memory β€” never runs its finally, so the broker waits out the session timeout; make sure orderly shutdowns never end that way.)
  5. GroupId correctness β€” is the value identical, byte-for-byte, across every instance meant to share partitions? A typo silently creates a second group instead of failing loudly.
  6. Exception handling β€” does the loop catch ConsumeException narrowly, check Error.Code, and send bad records to a dead-letter path instead of only logging them?
πŸ”§ Check🎯 What it verifies⚠️ Symptom if wrong
πŸ“ Commit configurationProgress recorded after the workSilent loss or duplicates on crash or restart
πŸ”’ One poll loopSingle instance, single threadOut-of-order processing, offsets committed past unfinished work
⏱️ Handler durationFast return to the poll loopMax poll interval exceeded, unwanted rebalances
πŸšͺ Shutdown handlingClose() on every exit pathBroker waits out session timeout to notice departure
🏷️ GroupId correctnessIdentical value across instancesEach instance reads the full topic instead of sharing it
🧯 Exception handlingNarrow catch, bad records dead-letteredFailed records vanish with no trace

⚠️ Common Mistake: Treating this as a one-time setup review rather than something to re-check after refactors. A handler fast at launch can quietly grow slow when a new dependency call is added inside the loop, reintroducing exactly the problem the checklist was meant to catch.

✍️ Exercise (harder): more throughput without more partitions

The payments topic has 6 partitions, keyed by customer ID, and the group already runs 6 instances. Each record's handler spends most of its time waiting on an external API, and the team needs roughly 4Γ— the throughput. Repartitioning is off the table this quarter, because it would remap keys to different partitions. Sketch a design that gets the throughput while keeping each customer's payments in order and keeping at-least-once. Say what the poll loop does, what the workers do, and what offset gets committed.

Check your answer

More instances cannot help: 6 partitions cap the group at 6 active members. The extra concurrency has to come from inside each instance, and it has to respect keys, not just partitions:

  • Poll loop: stays single-threaded. It hands each record to a worker chosen by hash(key) % N, so every record for one customer goes to the same worker and is processed in order. It pauses consumption (consumer.Pause) when the workers' bounded queues are full, and resumes when they drain.
  • Workers: process their own queue sequentially and report (partition, offset) when a record finishes.
  • Commits: for each partition, track which offsets have completed, and commit (or store) only the offset after the highest contiguous completed one. If offset 51 finishes before 50, the partition's bookmark waits at 50.
  • Rebalances: when partitions are revoked, stop handing out their records, wait for or abandon their in-flight work, and commit what is contiguous before letting go.
  • Duplicates: after a crash, everything past the last contiguous commit replays, so the payment call must be idempotent (for example, an idempotency key derived from the payment ID).

With N workers per instance you get up to 6 Γ— N concurrent calls, while each customer's payments stay in order.

What You Now Understand That You Didn't Before

At the start of this lesson, a single consumer reading a topic serially was the obvious approach and its throughput ceiling wasn't obviously a design problem. You now have the vocabulary to reason about why that ceiling exists, how a consumer group removes it up to the limit set by partition count, why an offset commit is a bookmark rather than a receipt, and why the .NET client's defaults quietly allow both loss and duplication unless you choose a commit pairing. Those ideas β€” ownership is partition-scoped, commits are checkpoints, and where the checkpoint lands is decided by your commit configuration, defaults included β€” are what every later lesson assumes you already have.

Where This Goes Next

  • KRaft & Modern Kafka covers the cluster side, including the newer consumer group protocol mentioned at the start of this lesson.
  • Raw Confluent.Kafka Client goes through the client API surface in more detail.
  • Consumer & Manual Commits goes deep on commit timing, the rebalance callbacks, cooperative versus eager rebalancing, and max.poll.interval.ms β€” the tools for choosing deliberately between at-least-once and at-most-once instead of getting whichever your configuration accidentally produces.
  • Delivery Semantics & Schema Strategy builds on the duplicate and loss windows traced here.

πŸ’‘ Pro Tip: As you move into those lessons, keep re-anchoring new detail back to the two recap ideas β€” partition-scoped ownership and offsets-as-checkpoints. Every mechanism you're about to learn is an elaboration of one of those, not a replacement.