KubeMQ
ConnectorsKafkaHow-to guides

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 next storage 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:

AcknowledgementEffect
AcceptTerminal — the record is consumed; the group's start offset advances past it once every earlier record is also resolved.
ReleaseThe record becomes available for redelivery (to this or another member), with its delivery count incremented.
RejectTerminal, like Accept, but signals "skip this record" rather than "processed successfully" — it is never delivered again.
RenewExtends 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-group

Supported clients

ClientStatus
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=earliest on the server. by_duration:<ISO-8601 duration> (for example by_duration:PT1H) is also accepted; it compares KubeMQ's own ingest time, not the producer's CreateTime.
  • To set it for one group, alter share.auto.offset.reset on 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-offsets before 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 it

A 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:

KeyDefaultAllowed rangeOutside the range
share.record.lock.duration.ms3000015000–60000refused INVALID_CONFIG
share.session.timeout.ms4500045000–60000refused INVALID_CONFIG
share.heartbeat.interval.ms50005000–15000refused INVALID_CONFIG
share.delivery.count.limit52–10clamped to the range, never refused
share.partition.max.record.locks2000100–4000clamped to the range, never refused
share.auto.offset.resetlatestlatest, earliest, by_duration:<ISO-8601>refused INVALID_CONFIG
share.isolation.levelread_uncommittedread_uncommitted, read_committedrefused INVALID_CONFIG
share.renew.acknowledge.enabletruetrue, falsefalse 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=3

A 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 pass max.poll.records, but always serves at least one batch. The default batch_optimized mode also counts stored batches, so a fetch can exceed max.poll.records when the producer batches.
  • read_committed groups. Set share.isolation.level=read_committed on 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 retriable COORDINATOR_NOT_AVAILABLE until unused share groups are deleted.

Managing share groups

  • List and describe: kafka-share-groups.sh --list, --describe --group <g> with --members or --offsets, or the admin client's ListGroups with a share type filter. A group lists as Stable while a member is attached and Empty once idle.
  • Delete: kafka-share-groups.sh --delete --group <g> (DeleteGroups) works once no consumer is attached; while one is, it answers NON_EMPTY_GROUP.
  • Move start offsets: kmq kafka share-groups reset-offsets <group> --topic <t> with --to-earliest, --to-latest or --to-offset N, then --execute. It is refused while a consumer is attached. DeleteShareGroupOffsets resets 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_flight and kubemq_kafka_share_group_redelivered, labelled group, topic, partition. In-flight and redelivered count stored batches, the same unit as the in-flight limit.
  • API: GET /api/kafka/share-groups and GET /api/kafka/share-groups/:id return 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 list and kmq 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:

  1. Stop share consumers and share-group configuration changes until every node runs the new release, then upgrade normally.
  2. Set CONNECTORS_KAFKA_SHARE_FEATURE_VERSION=1 on each node as you upgrade it, and raise it back to 2 (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.
  • ShareFetch answers 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.

Was this page helpful?

On this page