acks=0 never waits and never learns of a loss; acks=1 waits for the leader's log, lost if the leader dies before replicating; acks=all, the default, waits for every in-sync replica (The In-Sync Replica Set). The module below, reused in this section, measures throughput, batch size and delivery latency; as a script it compares the levels on the six-node cluster.
"""producer_lab.py: measure a producer configuration on BookNest's order events."""
import json, statistics, time
from confluent_kafka import Producer
EVENTS = [(str(json.loads(line)["order_id"]), line.rstrip())
for line in open("data/order_events.jsonl")]
def bench(topic, config, n=100_000, rate=None, bootstrap="localhost:33092"):
waits, failed, batches = [], [], [0, 0] # latencies, errors, [records, batches]
def stats(js): # librdkafka statistics, every second
b = json.loads(js)["topics"].get(topic, {}).get("batchcnt", {})
batches[0] += b.get("sum", 0); batches[1] += b.get("cnt", 0)
p = Producer({"bootstrap.servers": bootstrap, "statistics.interval.ms": 1000,
"queue.buffering.max.messages": 500_000, "stats_cb": stats, **config})
start = time.perf_counter()
for i, (key, value) in enumerate(EVENTS[:n]):
if rate: # pace the producer at `rate` msg/s
time.sleep(max(0, start + i / rate - time.perf_counter()))
p.produce(topic, value, key, on_delivery=lambda e, m:
failed.append(e) if e else waits.append(m.latency()))
p.poll(0) # serve delivery callbacks
p.flush(); secs = time.perf_counter() - start
time.sleep(1.1); p.poll(0) # one last statistics callback
q = statistics.quantiles(waits, n=100)
return (f"{n / secs:8,.0f} msg/s {batches[0] / max(batches[1], 1):5.0f}/batch "
f"p50 {q[49] * 1000:6.1f} ms p99 {q[98] * 1000:6.1f} ms {len(failed)} failed")
if __name__ == "__main__": # acks, six-node cluster, 2,000 msg/s
for acks in ("0", "1", "all"):
print(f"acks={acks:<4}", bench("booknest.acks", {"acks": acks}, n=10_000,
rate=2_000, bootstrap="localhost:33194"))Output
acks=0 2,000 msg/s 5/batch p50 3.6 ms p99 6.9 ms 0 failed acks=1 1,927 msg/s 5/batch p50 9.5 ms p99 164.1 ms 0 failed acks=all 1,486 msg/s 5/batch p50 1491.7 ms p99 1967.9 ms 0 failed
Across seven runs on this busy shared host the median was 3.6-4.5 ms for acks=0, 5-17 ms for acks=1 and 7.5 ms to 1.5 s for acks=all, which waits for the slowest follower. Keep acks=all for business events.