Share Groups

Share Groups and Acknowledgment Modes

A share consumer acknowledges each record as ACCEPT (done), RELEASE (deliver again, possibly to another consumer) or REJECT (never deliver again); since Kafka 4.2 129 Java clients can also RENEW a lock for a slow job (KIP-1222). In implicit mode, the default, the next poll accepts everything; in explicit mode the application acknowledges record by record. A record not acknowledged before share.record.lock.duration.ms (30 s) is released, and after share.delivery.count.limit (5) deliveries it is archived like a rejected one. BookNest turns 2,000 paid orders into packing jobs on a three-partition topic and runs four packers.

share_packers.py: four share consumers on three partitionsPython
"""share_packers.py: four packers share a 3-partition queue and acknowledge every job."""
import json, threading
from collections import Counter
from confluent_kafka import AcknowledgeType, ShareConsumer
stats = Counter()
def packer(name):
    c = ShareConsumer({"bootstrap.servers": "localhost:33092", "group.id": "booknest-packers",
                       "share.acknowledgement.mode": "explicit"})
    c.subscribe(["booknest.packing-jobs"])
    idle = 0
    while idle < 10:                                  # stop after 10 empty one-second polls
        msgs = c.poll(1.0)
        idle = 0 if msgs else idle + 1
        for m in msgs:
            order = json.loads(m.value())["order_id"]
            if order % 500 == 0:                      # bad address: never deliver again
                c.acknowledge(m, AcknowledgeType.REJECT); stats["rejected"] += 1
            elif order % 7 == 0 and m.delivery_count() == 1:   # printer jam: try again later
                c.acknowledge(m, AcknowledgeType.RELEASE); stats["released"] += 1
            else:
                c.acknowledge(m, AcknowledgeType.ACCEPT); stats[name] += 1
                stats["redelivered"] += m.delivery_count() > 1
        c.commit_sync()                               # send this batch's acknowledgements
    c.close()
threads = [threading.Thread(target=packer, args=(f"packer-{i}",)) for i in range(1, 5)]
for t in threads: t.start()
for t in threads: t.join()
print(dict(sorted(stats.items())))
Output
{'packer-1': 653, 'packer-2': 318, 'packer-3': 523, 'packer-4': 502, 'redelivered': 280, ...}
PARTITION LAG
0 0
1 0
2 0
...

All four packers worked, more than a consumer group could use. The 280 released jobs came back and were accepted on their second delivery, the four rejected ones were archived, and lag is zero. listings/l07610_run.sh creates the topic and sets share.auto.offset.reset=earliest on the group (new share groups start at the end).