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 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()))){'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).