One Slow Message. One Frozen Partition. Three Idle Workers. Kafka 4.2 finally lets many consumers share one partition — record by record. CONSUMER GROUP partition 0 ! everything behind it waits C1 stuck C2 idle C3 idle C4 idle 1 partition ⇒ at most 1 worker lag climbs · committed offset frozen SHARE GROUP partition 0 ! C1 slow C2 C3 C4 1 partition ⇒ many workers each record acked on its own KIP-932 "Queues for Kafka" — early access in 4.0, production-ready in 4.2

You have a Kafka topic called image-jobs with 6 partitions. Every message says "make thumbnails for this upload." On a normal day, 6 consumers keep up easily. Then a marketing campaign lands, uploads jump 10×, and lag starts to climb. So you do the obvious thing: scale the consumer deployment from 6 pods to 30. Lag keeps climbing. Twenty-four of your new pods are running, healthy, connected to Kafka — and doing absolutely nothing. Meanwhile one partition hasn't moved in four minutes, because the message at the front of it is a 400 MB TIFF that takes ages to resize, and nothing behind it can go until it's done.

Both of those problems — the workers you can't use and the one message that freezes everything behind it — come from the same design decision at the heart of Kafka: a partition belongs to exactly one consumer in a group, and progress on that partition is a single number. That rule is what makes Kafka great for ordered event streams. It is also why, for fifteen years, people who wanted a plain work queue either bent Kafka into shapes it didn't like or ran RabbitMQ or SQS next to it.

Kafka 4.2 (February 2026) ships the fix as production-ready: share groups, from KIP-932, "Queues for Kafka." This post tells the story in order: the problem, how we solved it earlier, and how Kafka solves it now — then puts earlier and now side by side, problem by problem, and ends with the code, the settings, and the guarantees you give up when you switch.

This post in three steps 1 · THE PROBLEM One owner per partition idle workers + head-of-line blocking 2 · EARLIER FIXES Workarounds in your code partitions, retry topics, Parallel Consumer, SQS 3 · NOW Share groups (Kafka 4.2) the broker tracks every single message

The feel: supermarket lanes vs. the bank line

Think of a supermarket. Each checkout lane has one cashier, and each customer picks a lane and stays in it. It works well until someone at the front of lane 3 needs a price check. Now lane 3 is frozen. The people behind that customer can't jump to another lane, and the cashiers in lanes 1, 2 and 4 might be standing there with nobody to serve. Adding more cashiers doesn't help either — a new cashier with no lane of their own just stands around.

Now think of a bank. There's one line, and whichever teller is free calls "next." If one customer takes twenty minutes, that teller is busy, but everyone else keeps moving through the other tellers. Add a teller and the line moves faster, straight away.

A classic Kafka consumer group is the supermarket: partitions are lanes, consumers are cashiers, and each lane has exactly one cashier. A share group turns each partition into the bank line: many consumers can take the next available record from the same partition, and a slow record only ties up the one consumer handling it.

Consumer groups give you order by giving each partition one owner. Share groups give you throughput by letting any free worker take the next record. You can't fully have both — so Kafka now lets you choose.
Part 1The problemWhy adding workers stops helping, and why one message can freeze a partition

How a classic consumer group works

Two rules explain almost everything about the old model.

Rule 1: one partition, one consumer. When consumers join a group, the group splits the topic's partitions between them, and each partition goes to exactly one member. With 6 partitions and 4 consumers, two consumers get two partitions each. With 6 partitions and 8 consumers, two consumers get nothing at all. They join, they heartbeat, they sit idle as hot spares.

Rule 2: progress is a single offset per partition. Records in a partition are numbered 0, 1, 2, 3… (the offset). The consumer group doesn't remember which records were processed. It remembers one number per partition — the committed offset — and that number means "everything below this is done." Commit 105 and you are saying offsets 0 through 104 are all handled.

Consumer group: parallelism is capped at the partition count topic: image-jobs P0 P1 P2 P3 consumer 1 ← P0 consumer 2 ← P1 consumer 3 ← P2 consumer 4 ← P3 consumer 5 · idle consumer 6 · idle each partition: exactly one owner, one committed offset Consumers 5 and 6 are paid for, connected, heartbeating — and get zero records.

This design is not a mistake. It gives you two things that are very hard to get any other way:

  • Ordering per key. All events for user-42 hash to the same partition, one consumer reads that partition, in order. "Account opened" is always processed before "account closed."
  • Tiny, cheap progress tracking. Millions of records per second, and the group only has to store one integer per partition. Replay is trivial: move the number back.

For event streams, CDC pipelines and stream processing, that is exactly the right deal. The trouble starts when the messages aren't events in a story — they're independent jobs: send this email, resize this image, call this webhook, run this report. Jobs don't care about order. They care about getting done, by whoever is free.

Problem 1: the parallelism ceiling

Because of Rule 1, the most consumers that can do useful work is the number of partitions. Six partitions, six workers, full stop. That turns partition count — a storage decision you made months ago — into a hard cap on how fast you can process today.

The usual answer is "just make more partitions." It works, but every option costs something:

  • Partitions aren't free. Every partition is a set of log segment files, index files and open file handles on every replica, plus a leader to elect and replication traffic to carry. Plenty of teams created 200 partitions "for headroom" on a topic that needed 12 for its data.
  • You can add partitions but never remove them. Over-provision for a Black Friday peak and you carry those partitions forever.
  • Adding partitions reshuffles keys. The default partitioner picks a partition from hash(key) % numPartitions. Change the count and user-42 starts landing in a different partition, so the per-key ordering you relied on breaks at the moment of the change.
  • You have to guess the peak ahead of time. Autoscaling consumers is easy. Autoscaling partitions is not really a thing.
What about KIP-848? Kafka 4.0 also made the new consumer rebalance protocol (KIP-848) generally available. It makes rebalances incremental and far less disruptive — a big improvement — but it doesn't change Rule 1. A partition still has at most one owner in a consumer group. The ceiling is the same; you just reach it more smoothly.

Problem 2: head-of-line blocking, in depth

Head-of-line blocking is a general name for one idea: the item at the front of a line is stuck, so everything behind it is stuck too — even items that could have been handled right away by someone else. The same term shows up in networking (one lost TCP packet delays every HTTP/2 stream behind it), in CPU pipelines and in switch buffers. In Kafka it comes straight from Rule 2.

Walk through it slowly. Your consumer polls a batch: offsets 100 to 110. Offset 101 is the 400 MB TIFF. The other ten are small JPEGs that take 50 ms each.

Why one slow record freezes a whole partition 100 101 102 103 104 105 106 107 108 109 110 done 4 min… 50 ms each — could be done, but they're waiting in line committed offset = 101 "everything below 101 is done" The offset can't jump the gap: • committing 111 would claim 101 is done — it isn't • so the pointer stays at 101 until the TIFF finishes • crash now ⇒ 101–110 are all delivered again Progress is one number, so it can only move as fast as the slowest record at the front.

If your consumer processes one record at a time — which is what the standard poll loop does — records 102 to 110 simply wait. Lag on that partition grows, and no other consumer can help, because no other consumer is allowed to read partition 0.

"Fine, I'll process the batch in parallel threads inside my consumer." Now 102 to 110 finish in a few hundred milliseconds. But you still can't commit them. The only thing you can commit is a single offset, and committing 111 would mean "101 is done too." So you're stuck with a choice between two bad options: commit early and lose 101 if you crash, or hold the commit and redo 102–110 if you crash. You also need to track which offsets are finished yourself, and poll carefully so memory doesn't blow up while 101 drags on.

It gets worse in production. Here is how head-of-line blocking really shows up on call:

  • The slow message. One job takes minutes instead of milliseconds. Everything behind it in that partition waits. The graph shows lag on one partition climbing while the others sit near zero.
  • The rebalance loop. If handling one poll's worth of records takes longer than max.poll.interval.ms (5 minutes by default), the group decides your consumer is dead and moves its partitions to someone else. The new owner starts from the last committed offset — which is still 101 — picks up the same slow message, and gets kicked out the same way. Your whole group spends its time rebalancing.
  • The poison pill. A malformed message makes your handler throw every time. If you retry forever, the partition is stuck for good. If your process crashes, it restarts, reads the same message from the committed offset, and crashes again. A single bad record takes out a slice of your pipeline.
  • The flaky dependency. The webhook for one customer times out after 30 seconds, every time. Every job for every other customer that happens to share that partition waits behind it.
The key point: in a consumer group, you cannot say "skip this one for now, I'll come back to it." There is no per-record state — just a line and a pointer. Every workaround people built is really a way to fake per-record state on top of a system that doesn't have it.
Part 2How we solved it earlierFive workarounds teams built on top of Kafka — and the bill each one left

The earlier fixes: five workarounds (and why each one hurt)

People didn't just sit and suffer. For about fifteen years, teams running job queues on Kafka built their own fixes. Almost every production Kafka setup you'll meet uses one or more of these five. Each one works, and each one leaves a bill.

Fix 1: "Just add more partitions"

The first thing everyone tries. Go from 6 partitions to 60, run 60 consumers, and each lane gets shorter. It does help with throughput. But we already saw the cost: broker overhead, no way to shrink later, and keys moving to new partitions. It also doesn't fix head-of-line blocking at all — it only makes each lane shorter. A stuck record still freezes everything behind it in its partition. You've just made the blast radius 1/60th of the topic instead of 1/6th.

Fix 2: "Catch it, log it, move on"

for (ConsumerRecord<String, String> record : records) {
    try {
        makeThumbnails(record.value());
    } catch (Exception e) {
        log.error("failed, skipping offset {}", record.offset(), e);   // 🙈
    }
}
consumer.commitSync();   // the failed job is now "done" forever

This keeps the partition moving, and it's what a lot of code quietly does. The price is silent data loss: a failed job counts as processed, and the only record that it ever existed is a log line nobody reads. It's fine for a view counter. It's not fine for "charge this customer" or "send this password reset email."

Fix 3: Retry topics and a dead-letter topic

This is the grown-up version, made popular by Uber's engineering blog in 2018 and built into Spring Kafka as @RetryableTopic. When a job fails, you don't retry it in place. You publish a copy to a retry topic, commit the original, and move on. A separate consumer reads the retry topic after a delay and tries again. After a few failed rounds, the record lands in a dead-letter topic (DLQ) for a human to look at.

The old way to retry: copy the message to another topic image-jobs main topic jobs-retry-1 wait 1 min jobs-retry-2 wait 10 min jobs-dlq manual review failfailfail consumer consumer consumer replay tool? 4 topics, 4 consumers, producer code in your handler — per job type, per team. It works. You also just built a message queue by hand.
// Spring Kafka hides the plumbing, but the topics are still real
@RetryableTopic(attempts = "4", backoff = @Backoff(delay = 1000, multiplier = 2.0))
@KafkaListener(topics = "image-jobs")
public void handle(String job) {
    makeThumbnails(job);            // throw ⇒ Spring re-publishes to a retry topic
}

@DltHandler
public void deadLetter(String job) {
    alerting.notify("thumbnail job gave up: " + job);
}
The clever trick inside this pattern. Why a separate topic per delay, instead of one retry topic? Because a Kafka consumer can only wait for the record at the front. If every record in jobs-retry-1 has the same 1-minute delay, they are already sorted by "time to run," so the consumer just waits for the front record and never blocks a record that's due sooner. Mix delays in one topic and you've rebuilt head-of-line blocking inside your retry system.

What it costs:

  • Topic sprawl. Every job type needs its own set of retry and DLQ topics, consumers, alerts and permissions. Multiply by every team.
  • Your handler becomes a producer. "Publish to retry, then commit the original" is two steps. Crash between them and you either lose the job or retry it twice.
  • Order is gone anyway. A retried job runs minutes after the jobs that came after it.
  • The DLQ is a graveyard. Someone has to build tooling to inspect and replay it, and usually nobody does.
  • It doesn't help slow jobs. Retry topics handle failures. A job that is slow but succeeding still blocks its partition.

Fix 4: Process in parallel inside the consumer

If one consumer owns the partition, make that one consumer do many things at once. Confluent's open-source Parallel Consumer library is the best-known version: it reads a partition with a normal consumer, runs records on many threads, and tracks which offsets finished out of order.

ParallelConsumerOptions<String, String> options = ParallelConsumerOptions.<String, String>builder()
        .ordering(ProcessingOrder.UNORDERED)   // or KEY: parallel across keys, ordered within a key
        .maxConcurrency(100)                   // 100 jobs in flight from one consumer
        .consumer(kafkaConsumer)
        .build();

ParallelStreamProcessor<String, String> processor =
        ParallelStreamProcessor.createEosStreamProcessor(options);
processor.subscribe(List.of("image-jobs"));
processor.poll(ctx -> makeThumbnails(ctx.getSingleConsumerRecord().value()));

How does it get around the single-offset rule? It still commits the lowest safe offset (101 in our example), but it also packs a compressed list of "these offsets after 101 are already done" into the commit metadata string. After a restart, it reads that list back and skips the finished ones. It's a smart trick — per-record state, smuggled into a field meant for a short note.

What it costs:

  • Still one owner per partition. You scale threads inside a pod, not pods. Your 24 idle pods are still idle; you just gave 6 pods 100 threads each.
  • A limit on how far ahead it can get. Commit metadata has a size limit, so a big enough gap behind a stuck record still forces the library to slow down.
  • A library, not the platform. It's Java-only, it's another dependency to understand, and all the state lives in the client. The broker has no idea any of this is happening, so the normal tools can't show you which jobs are in flight.

Fix 5: Put a real queue next to Kafka

Finally, the honest answer many teams gave: "Kafka isn't a queue, so let's use a queue." Events go to Kafka. Jobs go to RabbitMQ or Amazon SQS, sometimes copied over by a bridge service or a Kafka Connect sink. SQS in particular has exactly what we've been missing: any worker takes any message, a visibility timeout hides a message while someone works on it, and a maxReceiveCount redrive policy sends repeat failures to a DLQ.

What it costs: a second system to run, secure, monitor, pay for and be on call for. Data is copied between them, so there are two places a job can get lost and two sets of metrics to line up at 3 AM. And you lose what made Kafka attractive in the first place — a job queue you can replay from the log.

What all five have in common: each one is a way to track the state of individual messages — outside the broker, because the broker couldn't. Retry topics store state as extra copies of messages. The Parallel Consumer stores it in commit metadata. SQS stores it in a different product. Share groups move that state into Kafka itself.
Part 3The new solution: share groupsWhat Kafka built, when it shipped, and how it works inside the broker

When and why share groups arrived

KIP-932, titled "Queues for Kafka," was proposed in 2023 with a clear goal: give Kafka the semantics of a work queue — any free consumer takes the next message, each message is acknowledged on its own, failures are retried a limited number of times — on top of the same durable, replayable log. No second system and no copying data. The topic doesn't change at all; a share group is just a different way of reading it.

ReleaseDateStatus of share groups
Kafka 4.0March 2025Early access. Not for production; a cluster that used it could not be upgraded.
Kafka 4.1September 2025Preview. Rewritten and not compatible with 4.0 early access; enabled with a feature flag.
Kafka 4.2February 2026Production-ready. Adds the RENEW acknowledgement (KIP-1222), share-group lag metrics (KIP-1226) and a configurable acquire mode (KIP-1206).
Kafka 4.3May 2026Per-group overrides for delivery limit and in-flight window (KIP-1240).
Next—Built-in dead-letter queues (KIP-1191) are accepted and in development, but not in any release as of 4.3.1.

How share groups work

A share group looks like a consumer group from the outside: consumers use a group.id, subscribe to topics, and poll. Two things are fundamentally different.

1. Many consumers can read the same partition. Assignment happens on the broker (the built-in simple assignor), and it is allowed to give one partition to several members. You can run 30 consumers against a 6-partition topic and all 30 will get work. The cap on members is a broker setting (group.share.max.size, 200 by default), not your partition count.

2. The broker tracks every in-flight record. Instead of one committed offset per partition, the broker keeps a small state machine for each record that has been handed out. When a consumer fetches, the broker gives it a batch of records and marks them Acquired by that consumer, with a time-limited lock (30 seconds by default). No other consumer will get those records while the lock holds. The consumer then acknowledges each record on its own.

The life of one record in a share group Available waiting for a consumer Acquired locked to one consumer Acknowledged done — never again Archived given up — never again fetched deliveryCount + 1 RELEASE or lock timeout only if deliveryCount < limit ACCEPT REJECT or limit reached RENEW "still working" — extend lock (4.2+) Default limit is 5 deliveries — a poison pill is tried 5 times, then archived, and the partition moves on. Acknowledged and Archived are terminal: the record is "delivery complete" either way.

Look at what this state machine gives you, compared with the old list of problems:

  • Slow record? It stays Acquired by one consumer. The records behind it are Available and go to other consumers. Nobody waits in line.
  • Consumer crashed mid-job? Its lock expires, the record goes back to Available, and another consumer picks it up. No rebalance of the whole group is needed for that.
  • Transient failure? The consumer acknowledges with RELEASE, and the record goes straight back to Available for another try.
  • Poison pill? Every delivery increments a counter. Once it reaches group.share.delivery.count.limit (default 5), the record is Archived instead of redelivered. The poison pill costs five attempts, not your pipeline. Or your code can see the problem right away and send REJECT.

The in-flight window: how the broker keeps this cheap

Tracking state for every record in a topic that holds billions of records would be impossible. The trick is that the broker only tracks a window of records per partition — the ones currently in play. KIP-932 calls this a share-partition, bounded by two offsets:

  • SPSO (Share-Partition Start Offset): everything below this is finished (acknowledged or archived). It behaves like the old committed offset.
  • SPEO (Share-Partition End Offset): everything at or above this hasn't been handed out yet.

Only records between the two carry per-record state. When the records at the start of the window all reach a terminal state, SPSO slides forward. The size of the window is capped by group.share.partition.max.record.locks (2,000 by default in 4.2): once that many records are in flight on a partition, the broker stops handing out more until some finish.

Per-record state lives only inside the in-flight window 98 99 100 101 102 103 104 105 106 107 108 109 110 C1 · slow C2 retry poison C3 SPSO = 100 SPEO = 108 in-flight window (≤ 2,000 records by default) finished not delivered yet Acknowledged Acquired (locked) Available again Archived 101, 102 and 105 are done even though 100 isn't — the gap no longer blocks anyone.

Notice that SPSO is still stuck at 100 because of the slow record — the "start of the window" still can't pass an unfinished record. The difference is that it no longer matters for throughput. Records 101 to 107 are being processed and acknowledged by other consumers. The only thing the slow record holds back is the window's start; it only becomes a problem if thousands of records pile up behind it and the window fills. That's a much softer failure than "the whole partition stops."

Under the hood: where the state lives

Per-record state has to survive broker restarts and leader changes, or a failover would redeliver everything in flight. Kafka splits the work between two server components.

Who does what inside the cluster consumer A consumer B consumer C Group coordinator membership + assignment (many members per partition) Partition leader holds the share-partition in memory: SPSO, SPEO, per-record state, locks, delivery counts hands out records + applies acks Share coordinator durable copy of the state __share_group_state ShareSnapshot + ShareUpdate ShareGroupHeartbeat ShareFetch (acks can ride along) ShareAcknowledge Write…State Read…State On leader failover, the new leader reads the state back from the share coordinator and carries on.
  • The group coordinator handles membership. Consumers send ShareGroupHeartbeat requests, and the coordinator returns their assignment. Share groups only use this new heartbeat protocol, built on the same foundation as KIP-848.
  • The partition leader is where the action is. It holds each share-partition's window in memory, answers ShareFetch requests by acquiring records for the caller, and applies acknowledgements. Acks can be sent on the next ShareFetch (piggybacked, so no extra round trip) or on their own with ShareAcknowledge.
  • The share coordinator makes the state durable. When the window changes, the leader writes it via WriteShareGroupState to the internal topic __share_group_state — a full ShareSnapshot now and then, and smaller ShareUpdate records in between, the same snapshot-plus-delta idea as a database write-ahead log. On failover the new leader calls ReadShareGroupState and resumes with the same window.
Why state is kept per batch, not per message. Records are acquired and tracked in batches wherever possible, and split only when different records in a batch end up in different states. That's how the broker keeps the bookkeeping close to the cost of a normal fetch. It's also why Kafka 4.2 added share.acquire.mode: the default batch_optimized lines up with batch boundaries and may return more than max.poll.records, while record_limit strictly respects it.
Part 4Earlier vs. now, problem by problemEvery old problem, the old fix, and what the broker does for you today

The big picture: who does the work

Step back and the whole change fits in one picture. Earlier, your application carried the bookkeeping: extra partitions, thread pools, offset tracking, retry producers, retry and dead-letter topics — or a whole second messaging system. Now the broker carries it, and the workers get simple again.

BEFORE your code does the bookkeeping image-jobs · 200 partitions your consumer app thread pool offset tracker retry producer rebalance hacks DLQ alerts retry-1 retry-2 dlq + a consumer for each one …or a bridge to RabbitMQ / SQS every team rebuilds this AFTER (share groups) the broker does the bookkeeping image-jobs · 6 partitions Kafka broker (share-partition) record locks delivery count per-record acks auto retry durable state 30 simple workers: poll → work → ack built into Kafka, same topic, same cluster

Problem by problem

Each card below reads left to right: the problem, how we fixed it earlier, and what share groups do now.

1 Not enough workers THE PROBLEM P0P1C1C2C3 idleC4 idle 2 partitions ⇒ only 2 workers extra consumers sit idle EARLIER FIX more lanes…can't removethem later Add partitions (forever) or add threads inside one pod NOW: SHARE GROUPS P0C1C2C3C4 Just add consumers — many share one partition 2 One slow message blocks the rest THE PROBLEM !4 min…waiting in line Everything behind it waits, even when workers are free EARLIER FIX !commit stuck here Parallel Consumer: done offsets hidden in commit metadata NOW: SHARE GROUPS C1slow…C2C3C2C4🔒 Slow message locked to one worker; the rest keep moving 3 Retrying a failed message THE PROBLEM ✗retry forever ⇒ stuckskip it ⇒ data lost No way to say "try this one again later" EARLIER FIX mainretry-1retry-2dlq+ a consumer for each topic Copy it to retry topics, run more consumers NOW: SHARE GROUPS AvailableAcquiredfetch · try 2 of 5RELEASE Broker puts it back and counts the attempt 4 A poison pill THE PROBLEM read same msg💥 crashrestart Crash → restart → same message → crash. Forever. EARLIER FIX jobs-dlq topic🔔 alerts you wire up yourself Dead-letter topic you build, watch and replay yourself NOW: SHARE GROUPS 12345Archiveddelivery count limit = 5 After 5 tries it's Archived (real DLQ coming: KIP-1191) 5 A worker crashes mid-job THE PROBLEM C1 ✗whole group rebalancesreplay since last commit Everyone pauses, and work since the last commit repeats EARLIER FIX tune session / poll timeoutsmake handlers idempotentaccept the replays Tune, make it safe to redo, and live with the replays NOW: SHARE GROUPS 30sAvailableC2 ✓no group-wide replay Lock expires; only its messages go back for another worker 6 Needing a real queue at all THE PROBLEM Kafkalog only — no per-message acks No per-message acks, retries or locks EARLIER FIX KafkaRabbitMQ/ SQSbridge2 systems to run and watch Run a second system and copy data across NOW: SHARE GROUPS image-jobs topicconsumer groupshare group One topic, one cluster: stream + queue side by side
If you know SQS, you already know share groups. The acquisition lock is the visibility timeout (both default to 30 seconds). The delivery count limit is maxReceiveCount. ACCEPT is DeleteMessage. RELEASE is setting the visibility timeout to zero. RENEW is ChangeMessageVisibility. The difference is that the messages still live in a Kafka log — retained, replayable, and readable by ordinary consumer groups at the same time.
In one sentence: before, Kafka only remembered where each group was in a partition. Now it can also remember what happened to each message.
Part 5Using it for realThe code, the settings, what you give up, and when to pick which

The code

Turning it on

Share groups are controlled by the share.version feature. On a new Kafka 4.2+ cluster it's enabled by default. On an upgraded cluster (or on 4.1, where it's a preview), all brokers must be on the new version, and then you turn it on:

bin/kafka-features.sh --bootstrap-server localhost:9092 \
  upgrade --feature share.version=1
Small test clusters: the internal __share_group_state topic defaults to replication factor 3 with min.isr 2. On a 1- or 2-broker dev cluster, set share.coordinator.state.topic.replication.factor and share.coordinator.state.topic.min.isr to 1 first, or share groups won't work.

A worker with explicit acknowledgements

The client is KafkaShareConsumer. It looks almost like a normal consumer on purpose — but there's no assign() and no seek(), because the broker decides who gets what.

Properties props = new Properties();
props.put("bootstrap.servers", "localhost:9092");
props.put("group.id", "thumbnail-workers");             // this is now a SHARE group
props.put("key.deserializer", StringDeserializer.class.getName());
props.put("value.deserializer", StringDeserializer.class.getName());
props.put("share.acknowledgement.mode", "explicit");    // default is "implicit"

try (KafkaShareConsumer<String, String> consumer = new KafkaShareConsumer<>(props)) {
    consumer.subscribe(List.of("image-jobs"));

    while (running) {
        ConsumerRecords<String, String> records = consumer.poll(Duration.ofMillis(500));

        for (ConsumerRecord<String, String> record : records) {
            try {
                makeThumbnails(record.value());
                consumer.acknowledge(record, AcknowledgeType.ACCEPT);   // done
            } catch (TransientException e) {
                consumer.acknowledge(record, AcknowledgeType.RELEASE);  // try again, maybe elsewhere
            } catch (BadInputException e) {
                consumer.acknowledge(record, AcknowledgeType.REJECT);   // never retry this one
            }
        }

        // Send the acks now and see per-partition results.
        // Without this, acks are sent with the next poll().
        Map<TopicIdPartition, Optional<KafkaException>> result = consumer.commitSync();
    }
}

A few things worth knowing about this loop:

  • Implicit vs explicit mode. In the default implicit mode you don't call acknowledge at all: calling poll() again (or commitSync/commitAsync) accepts every record from the previous poll. That's the simplest migration from a normal consumer. In explicit mode you acknowledge each record yourself, and you should acknowledge everything from one poll before calling poll() again.
  • Acks are batched. acknowledge() just records your decision locally. It reaches the broker on the next poll(), commitSync() or commitAsync(). If you use commitAsync, register setAcknowledgementCommitCallback to find out whether acks landed.
  • You can see how many times a record was tried. Records from a share consumer carry a delivery count, so a handler can log or alert when it's on attempt 4 of 5.
  • Long jobs: RENEW. If a job may run longer than the lock (30 s by default), acknowledge it with AcknowledgeType.RENEW and commit to extend the lock while you keep working. Then send ACCEPT when you finish. Before 4.2 the only option was a longer lock duration for the whole group.

Operating a share group

# Who is in the group, and how far along is each partition?
bin/kafka-share-groups.sh --bootstrap-server localhost:9092 \
  --describe --group thumbnail-workers              # --members, --offsets, --state

# A brand-new share group starts at the END of the topic by default.
# To process what's already there, change the group config:
bin/kafka-configs.sh --bootstrap-server localhost:9092 --alter \
  --entity-type groups --entity-name thumbnail-workers \
  --add-config share.auto.offset.reset=earliest

# Rewind (the group must have no active members)
bin/kafka-share-groups.sh --bootstrap-server localhost:9092 \
  --reset-offsets --group thumbnail-workers --topic image-jobs \
  --to-earliest --execute
The gotcha everyone hits first. share.auto.offset.reset defaults to latest. Start a new share group on a topic that already has a million jobs in it and you'll process none of them — only jobs produced after the group started. Unlike a normal consumer, this is a group config (set with kafka-configs), not a client property.

The settings that matter (Kafka 4.2 defaults)

SettingWhereDefaultWhat it controls
group.share.record.lock.duration.msbroker30000How long a consumer holds a record before it's handed to someone else. Per-group override: share.record.lock.duration.ms.
group.share.delivery.count.limitbroker5 (range 2–10)Attempts before a record is archived. Per-group override from 4.3: share.delivery.count.limit.
group.share.partition.max.record.locksbroker2000Size limit of the in-flight window per partition. Per-group override from 4.3.
group.share.max.sizebroker200Max members in one share group.
share.auto.offset.resetgrouplatestWhere a new group starts: latest, earliest or by_duration:<ISO-8601>.
share.isolation.levelgroupread_uncommittedSet read_committed to skip aborted transactional records.
share.acknowledgement.modeconsumerimplicitAuto-accept on next poll vs. per-record acks.
share.acquire.modeconsumerbatch_optimizedAlign with batches, or strictly respect max.poll.records (record_limit).

What you give up

Share groups aren't "consumer groups, but better." They're a different deal, and you should know exactly what you're trading.

  • No ordering guarantee. Records from one partition are handed to many consumers and finish in any order. A released record comes back after records that were behind it. Records inside a single delivered batch arrive in offset order, but across the partition there's no promise. If order-created must be processed before order-shipped, a share group is the wrong tool.
  • At-least-once, and only at-least-once. A lock can expire while your handler is still working — then another consumer gets the same record, and you process it twice. Share groups don't support exactly-once or transactional consumption. Handlers must be idempotent (an idempotency key, an upsert, a "processed" table). See the exactly-once post for why this was always true anyway.
  • Lock duration is a real trade-off. Too short and slow jobs get duplicated; too long and a crashed consumer's records sit locked until the timer runs out. Measure your p99 job time and set the lock above it — or use RENEW for long jobs.
  • Archived is not a dead-letter queue — yet. After the delivery limit, the record is simply marked Archived and skipped. It's still in the topic log, but nothing tells you about it or puts it anywhere you can inspect. Built-in DLQ support (KIP-1191) is accepted but not released as of 4.3.1. Until then, if a failed job matters, publish it to your own error topic from the handler before you REJECT it.
  • No per-key state or stream processing. Kafka Streams, joins, windowed aggregations — anything that assumes one consumer sees all records for a key — needs a classic consumer group.
  • No manual control. No assign(), no seek(), no static membership. You can reset a group's start point with the CLI, but a consumer can't jump around in the log on its own.
  • The window can still fill. A record stuck at the start holds SPSO back. If 2,000 records pile up behind it, the partition stops handing out new work until the lock expires. Head-of-line blocking is much softer, not gone.
  • It's new. Production-ready since February 2026, with an important deadlock fix in 4.2.1. Lag metrics only exist from 4.2. Run the latest patch release, and check that your client libraries (not just the Java client) support share groups before you plan around them.

Which one should you use?

Your workloadUseWhy
Events that tell a story per key (orders, accounts, CDC)Consumer groupPer-key ordering is the whole point
Stream processing, joins, aggregationsConsumer group / Kafka StreamsNeeds all records for a key in one place
Independent jobs: emails, thumbnails, webhooks, reportsShare groupScale workers freely; slow jobs don't block fast ones
Jobs with very uneven processing timesShare groupThis is exactly the head-of-line problem it solves
Need delayed delivery, priorities, per-message TTL todaySQS / RabbitMQShare groups don't offer these features

And remember you don't have to choose per topic. The same topic can be read by a consumer group and a share group at the same time. Your analytics pipeline reads image-jobs in order with a consumer group, while 30 thumbnail workers chew through it as a queue with a share group. One copy of the data, two ways to read it.

Share group
  • Consumers can outnumber partitions — scale workers, not partitions.
  • A slow or failing record only blocks itself.
  • Built-in retry with a delivery limit; poison pills get archived.
  • Crashed worker's records return after the lock expires.
  • No second messaging system to run.
What it costs
  • No ordering across the partition.
  • At-least-once only; duplicates are expected, not rare.
  • No built-in DLQ yet; archived records are silent.
  • Lock timing needs tuning to your job times.
  • New: fewer client libraries and less production history.

The takeaway

For fifteen years, Kafka had one way to read a topic, and it was built for ordered streams: one owner per partition, one number for progress. That's why adding consumers stopped helping after the partition count, and why one slow or broken message could freeze everything behind it. Three things worth carrying out of this:

  1. Head-of-line blocking comes from how progress is tracked. If progress is a single pointer, the pointer can only move past records that are finished, so the slowest record at the front sets the pace. Parallel threads don't fix that. Tracking state per record does.
  2. Order and parallelism pull against each other. A consumer group picks order; a share group picks parallelism. Decide what your data actually needs. Jobs rarely need order. Events usually do.
  3. Queue semantics mean at-least-once. Once any worker can take any record and locks can expire, duplicates are part of normal operation. Make every handler idempotent before you switch, and have a plan for archived records until built-in DLQs ship.

If you're running a job queue on Kafka today with 200 partitions "for headroom," retry topics you wrote yourself, or a RabbitMQ cluster that exists only because Kafka couldn't do queues, it's worth a test on 4.2. Point a share group at the existing topic, keep your consumer group running next to it, and watch 30 workers finally do 30 workers' worth of work.

Further reading. KIP-932 "Queues for Kafka" on the Apache Kafka wiki — the full design, including the record state machine and share-partition persistence; the Kafka 4.2 release announcement and upgrade notes (production readiness, KIP-1222 RENEW, KIP-1226 lag metrics, KIP-1206 acquire mode); KIP-1240 (per-group share configs in 4.3); KIP-1191 (dead-letter queues for share groups, in progress); the kafka-share-groups.sh section of the Kafka operations docs. Related here: the myth of exactly-once delivery (why idempotent handlers are non-negotiable), backpressure and flow control (what the in-flight window is really doing), and long-tail latency (why one slow item dominates a queue).