Share Groups
Queue-style consumption with Kafka share groups (KIP-932) on KubeMQ — per-record acquire and acknowledge, admin tools, per-group settings and upgrades.
Kafka share groups (KIP-932) give the Kafka connector a second, queue-style way to consume a topic. Instead of assigning whole partitions to consumers, individual records are acquired, processed, and acknowledged one at a time — closer to how a KubeMQ Queue behaves than to a classic consumer group.
Share groups are generally available (GA). The data plane — ShareGroupHeartbeat(76),
ShareFetch(78), ShareAcknowledge(79) — and the admin keys ShareGroupDescribe(77),
DescribeShareGroupOffsets(90), AlterShareGroupOffsets(91), and DeleteShareGroupOffsets(92)
are implemented and advertised. Share groups are certified with Java's KafkaShareConsumer
(kafka-clients 4.3) and the Kafka 4.3 console tools, on a single node and on three-node
clusters. The cluster runs included killing the share coordinator mid-run, with every record
checked afterwards: nothing lost, nothing accepted twice.
Soak runs of mixed Accept, Release, Reject and poison traffic showed no loss, no growth in memory
or open files, a bounded number of records in flight, and exactly one archive per poison record.
A few behaviors deliberately differ from Apache Kafka 4.3 — read
the share-group differences
before you move a production workload.
Prerequisites
- A KubeMQ server with the Kafka connector on (the default). Share groups run on the
nextstorage engine, which a fresh deployment selects automatically when Kafka is enabled. - A client with a share-consumer API — see Supported clients.
- On a cluster, every node on the same release before share traffic starts — see Upgrading a cluster.
How share groups differ from classic consumer groups
A classic consumer group assigns whole
partitions to members via Join/Sync/Heartbeat — one partition, one owning consumer at a time,
with a rebalance whenever membership changes. A share group throws that model out: every member
can receive records from every partition, and the unit of ownership is a single record (or a
contiguous batch), not a partition. There's no partition-assignment protocol to reason about —
just heartbeat-based membership (ShareGroupHeartbeat) plus per-record acquisition on fetch.
That makes a share group behave much more like a KubeMQ Queue: multiple workers pull from a shared backlog, a record goes to exactly one worker at a time, and a worker that fails to process it releases the record back for someone else to pick up — see the acquire/acknowledge cycle below.
Acquire, deliver, acknowledge
Where a classic consumer commits offsets in bulk, a share consumer acknowledges per record (or per contiguous batch) with one of these outcomes:
| Acknowledgement | Effect |
|---|---|
| Accept | Terminal — the record is consumed; the group's start offset advances past it once every earlier record is also resolved. |
| Release | The record becomes available for redelivery (to this or another member), with its delivery count incremented. |
| Reject | Terminal, like Accept, but signals "skip this record" rather than "processed successfully" — it is never delivered again. |
| Renew | Extends the acquisition lock on a record the consumer still holds, without completing it. Needs explicit acknowledgement mode; a group can turn it off with share.renew.acknowledge.enable=false. |
A record the client never acknowledges is released automatically once its acquisition lock
expires — 30 seconds by default, set per group with share.record.lock.duration.ms. That's the
same outcome as an explicit Release, and it also counts as a delivery attempt.
Redelivery has a limit. A record whose delivery fails 5 times (the default, set per group
with share.delivery.count.limit) is archived: it is never delivered to that group again,
so a single poison record can never wedge the partition for everyone else. The record's bytes
are untouched; a plain Fetch consumer on the same topic still sees it. A Release, a lock
expiry, a consumer leaving or being evicted, and a partition's leader changing (a killed node, a
restart, or each step of a rolling upgrade) each count as one failed attempt. Keep the limit at 5
or higher on groups that must ride through rolling restarts.
KubeMQ stores the records a producer sends in one batch together, and a share group acquires, releases and counts them together. A Release or a renew therefore acts on the whole stored batch: records in that batch you already accepted come back as duplicates. If your workload releases individual records often, have the producer send smaller batches.
Produce and share-consume
Both client libraries below drive the full flow — produce onto a plain topic, then acquire, process, and acknowledge from a share group. A new share group starts at the latest offset, so start the consumer before you produce, or set the group's start position first (see Start position).
package main
import (
"context"
"fmt"
"log"
"github.com/twmb/franz-go/pkg/kgo"
)
func main() {
ctx := context.Background()
// A share group has no partition assignment — every member shares the
// same pool of records, each one acquired individually.
sc, err := kgo.NewClient(
kgo.SeedBrokers("localhost:9092"),
kgo.ShareGroup("orders-share-group"),
kgo.ConsumeTopics("orders"),
)
if err != nil {
log.Fatal(err)
}
defer sc.Close()
for {
fetches := sc.PollFetches(ctx)
if errs := fetches.Errors(); len(errs) > 0 {
log.Fatal(errs[0].Err)
}
var accepted []*kgo.Record
fetches.EachRecord(func(r *kgo.Record) {
fmt.Printf("acquired offset=%d attempt=%d: %s\n", r.Offset, r.DeliveryCount(), r.Value)
// process the record here — on failure, MarkAcks with AckRelease
// or AckReject instead of AckAccept below.
accepted = append(accepted, r)
})
// franz-go accepts any record left unmarked on the next poll, so
// mark every record before polling again.
sc.MarkAcks(kgo.AckAccept, accepted...)
if err := sc.FlushAcks(ctx); err != nil {
log.Printf("flush acknowledgements: %v", err)
}
}
}Properties props = new Properties();
props.put("bootstrap.servers", "localhost:9092");
props.put("group.id", "orders-share-group");
// acknowledge() needs explicit mode; the default (implicit) accepts every
// record on the next poll.
props.put("share.acknowledgement.mode", "explicit");
props.put("key.deserializer", "org.apache.kafka.common.serialization.StringDeserializer");
props.put("value.deserializer", "org.apache.kafka.common.serialization.StringDeserializer");
try (KafkaShareConsumer<String, String> consumer = new KafkaShareConsumer<>(props)) {
consumer.subscribe(Collections.singleton("orders"));
while (true) {
ConsumerRecords<String, String> records = consumer.poll(Duration.ofMillis(5000));
for (ConsumerRecord<String, String> record : records) {
// process the record here — on failure, acknowledge RELEASE or
// REJECT instead of ACCEPT below.
consumer.acknowledge(record, AcknowledgeType.ACCEPT);
}
consumer.commitSync(); // sends the pending acknowledgements
}
}The Kafka 4.3 console tools work unchanged:
kafka-console-share-consumer.sh --bootstrap-server localhost:9092 --topic orders --group orders-share-group
kafka-share-groups.sh --bootstrap-server localhost:9092 --list
kafka-share-groups.sh --bootstrap-server localhost:9092 --describe --group orders-share-group --offsets
kafka-share-groups.sh --bootstrap-server localhost:9092 --delete --group orders-share-groupSupported clients
| Client | Status |
|---|---|
Java KafkaShareConsumer (kafka-clients 4.3) | Certified on a single node and on three-node clusters, including a killed share coordinator. |
Kafka 4.3 console tools (kafka-console-share-consumer.sh, kafka-share-groups.sh, kafka-configs.sh --entity-type groups) | Certified: consume, list, describe (members and offsets), delete, and per-group configuration. |
franz-go (kgo.ShareGroup) | Supported: the client the connector's own share-group tests and soak runs use. |
librdkafka share consumer (for example confluent-kafka Python's ShareConsumer) | Upstream preview. It passes an advisory single-node check; treat it as preview until librdkafka ships it as stable. |
kcat, kafkajs, Confluent.Kafka (C#), and the Ruby and Rust rdkafka bindings do not
expose a share consumer yet; use a classic consumer group with those clients.
Start position
A share group's start position on a partition is decided once, the first time the group
touches that partition; after that the group keeps its position. The default is latest,
matching Apache Kafka: records already on the topic when the group is created are not delivered
to it. Earlier releases started every new group at the earliest retained record.
- To start every new group at the beginning, set
CONNECTORS_KAFKA_SHARE_AUTO_OFFSET_RESET=earlieston the server.by_duration:<ISO-8601 duration>(for exampleby_duration:PT1H) is also accepted; it compares KubeMQ's own ingest time, not the producer'sCreateTime. - To set it for one group, alter
share.auto.offset.reseton that group (see Per-group configuration) before its first consumer joins. - To place a group at an exact offset, run
kmq kafka share-groups reset-offsetsbefore the first consumer joins. KubeMQ creates the group's start offset where it has none, so the first consumer starts there:
kmq kafka share-groups reset-offsets orders-share-group --topic orders --to-earliest # preview only
kmq kafka share-groups reset-offsets orders-share-group --topic orders --to-earliest --execute # write itA partition added to a topic after the group joined starts at the group's start position on the member's next heartbeat.
Per-group configuration
Every share group accepts kafka-configs.sh --entity-type groups --entity-name <group> for
these keys:
| Key | Default | Allowed range | Outside the range |
|---|---|---|---|
share.record.lock.duration.ms | 30000 | 15000–60000 | refused INVALID_CONFIG |
share.session.timeout.ms | 45000 | 45000–60000 | refused INVALID_CONFIG |
share.heartbeat.interval.ms | 5000 | 5000–15000 | refused INVALID_CONFIG |
share.delivery.count.limit | 5 | 2–10 | clamped to the range, never refused |
share.partition.max.record.locks | 2000 | 100–4000 | clamped to the range, never refused |
share.auto.offset.reset | latest | latest, earliest, by_duration:<ISO-8601> | refused INVALID_CONFIG |
share.isolation.level | read_uncommitted | read_uncommitted, read_committed | refused INVALID_CONFIG |
share.renew.acknowledge.enable | true | true, false | false refuses every renew with INVALID_RECORD_STATE |
kafka-configs.sh --bootstrap-server localhost:9092 --entity-type groups \
--entity-name orders-share-group --alter \
--add-config share.record.lock.duration.ms=20000,share.delivery.count.limit=3A change takes effect within one heartbeat interval of every live member, except
share.auto.offset.reset, which only applies to partitions the group has not started on yet.
The server-wide defaults and the ranges come from the CONNECTORS_KAFKA_GROUP_SHARE_* settings
(see the Kafka settings reference).
share.assignment.interval.ms is refused, because every member is always assigned every
partition. A group's configuration is kept when the group is deleted, so a group re-created
with the same name inherits it.
Limits and delivery behavior
- In-flight limit. Each share partition holds at most 2000 acquired or released but not yet
acknowledged stored batches by default (
share.partition.max.record.locks). The limit counts stored batches, not records, so with a batching producer more than 2000 records can be in flight. A Release does not free a slot; an Accept, a Reject or an archive does. At the limit a fetch returns no new records from that partition until one frees — not an error. - Acquire mode.
share.acquire.mode=record_limit(Kafka 4.2 and later clients) is honored in whole stored batches: KubeMQ stops before the batch that would passmax.poll.records, but always serves at least one batch. The defaultbatch_optimizedmode also counts stored batches, so a fetch can exceedmax.poll.recordswhen the producer batches. read_committedgroups. Setshare.isolation.level=read_committedon a group and it never reads past the partition's last stable offset, never receives a record from an aborted transaction, and never receives a transaction marker.- Group count. Share groups count against
CONNECTORS_KAFKA_MAX_GROUPS(default 10000), separately from classic groups. A share group keeps counting until it is deleted; once the cap is reached a new group is refused with the retriableCOORDINATOR_NOT_AVAILABLEuntil unused share groups are deleted.
Managing share groups
- List and describe:
kafka-share-groups.sh --list,--describe --group <g>with--membersor--offsets, or the admin client'sListGroupswith asharetype filter. A group lists asStablewhile a member is attached andEmptyonce idle. - Delete:
kafka-share-groups.sh --delete --group <g>(DeleteGroups) works once no consumer is attached; while one is, it answersNON_EMPTY_GROUP. - Move start offsets:
kmq kafka share-groups reset-offsets <group> --topic <t>with--to-earliest,--to-latestor--to-offset N, then--execute. It is refused while a consumer is attached.DeleteShareGroupOffsetsresets a group's start offset to the log start. - One id, one type: a group id in use as a share group cannot be used as a classic group, or
the other way round; the request answers
GROUP_ID_NOT_FOUND.
Monitoring
The share coordinator samples share groups every 5 seconds — up to 200 groups per round,
set by CONNECTORS_KAFKA_SHARE_METRICS_MAX_GROUPS_PER_TICK, taken in turn. A change reaches
every view within about 10 seconds.
- Prometheus:
kubemq_kafka_share_group_lag,kubemq_kafka_share_group_in_flightandkubemq_kafka_share_group_redelivered, labelledgroup,topic,partition. In-flight and redelivered count stored batches, the same unit as the in-flight limit. - API:
GET /api/kafka/share-groupsandGET /api/kafka/share-groups/:idreturn the same view from any node — members, and per partition the start offset, lag, in-flight, redelivered and highest delivery count, plus when the group was last sampled. - Dashboard: the Kafka page's consumer-groups tab lists share groups beside classic groups, with a Classic/Share filter and a detail page per share group.
- CLI:
kmq kafka share-groups listandkmq kafka share-groups describe <group>print the same view.
Upgrading a cluster
Upgrade every node before share traffic resumes. Share groups use internal replication commands that a node running an earlier release cannot apply. Such a node stops as soon as one is written, and keeps stopping until it is upgraded — on a three-node cluster that can cost quorum. Before a rolling upgrade, do one of these:
- Stop share consumers and share-group configuration changes until every node runs the new release, then upgrade normally.
- Set
CONNECTORS_KAFKA_SHARE_FEATURE_VERSION=1on each node as you upgrade it, and raise it back to2(a restart) once every node is upgraded.
At 1, the in-flight limit is not enforced, deleting a share group answers
UNSUPPORTED_VERSION, per-group configuration writes are refused, renew is refused,
record_limit is unavailable, read_committed groups are refused, and kafka-share-groups.sh --list does not work. Downgrading is supported only if the setting stayed at 1 the whole
time. Do not set share.isolation.level=read_committed on any group until every node runs the
new release.
Differences from Apache Kafka
The behaviors most likely to matter in a migration:
- A consumer that dies holding records gives them back only when their acquisition lock expires (30 seconds by default), not when its connection drops.
ShareFetchanswers as soon as it has a result instead of waiting up to the client's maximum wait.- The in-flight limit counts stored batches, not records.
The full list is in What Differs from Kafka.
Related
Consumer Groups
The classic Join/Sync/Heartbeat protocol, durable per-group offsets, and static membership — what a share group replaces the partition-assignment model with.
What Differs from Kafka
Every deliberate share-group difference from Apache Kafka 4.3, in one list.
Capabilities
Every implemented Kafka API and its version range, including the share-group keys.
Fitness matrix
What's drop-in, what works with a caveat, what's on the roadmap, and what's unsupported when running Kafka workloads on KubeMQ.
Was this page helpful?
Transactions & EOS
Exactly-once semantics on KubeMQ — the transactional producer, read_committed isolation, consume-transform-produce, and producer fencing.
Authentication
Authenticate Kafka clients to KubeMQ — SASL/PLAIN and SCRAM, OAUTHBEARER/OIDC federated tokens, mTLS client certificates, and the ACL authorization model.