Skip to content

Kafka clients

Kafka clients producing and consuming through the facade inside Queen, beside Queen's own clients and Kafka 4.3.1 on the same machines: rate, partition count, CPU, what made the facade fast, and what is still open.

Updated View as Markdown

Queen speaks the Kafka wire protocol from inside the broker process (QUEEN_KAFKA_EMBEDDED=true), so an existing producer or consumer can point at Queen and keep its code. We measured what that costs with the same load generator, built on franz-go, that drove Kafka itself. After a night of work on the facade, Kafka clients see the latency Queen’s own clients see, and at a million messages a second a little better.

The shape

The rig and load of the other pages: one topic of 200 partitions unless stated, 256-byte messages produced 100 at a time to one partition, 60 seconds per point, figures over the last 30. The clients were nine franz-go processes on three load machines, with acks=all, idempotent producers, lz4, a 5 ms linger and classic consumer groups. The Queen cluster ran the facade on every node, with the broker settings of the native runs. The builds are those of the facade work that was merged as b51e4419 (2.0.0-beta.5): qk3 on 2026-10-01, and the final build on 2026-10-02 for the row marked so. Runs: benchmark-queen/2026-09-30-kafka-pulsar/runs/queen-kafka/ (README-c3.txt there maps tags to builds).

Rate

End-to-end latency, p50 / p99 in milliseconds, one topic of 200 partitions. “Behind” means consumption under 97% of the offered rate, or a p99 over 5 seconds.

Offered msg/s Kafka clients through Queen Queen’s own clients Kafka 4.3.1
300,000 10 / 19 11 / 22 8 / 10
600,000 12 / 38 19 / 98 8 / 13
1,000,000 44 / 108; final build 35 / 91 61 / 146 8 / 18
1,500,000 185 / 3,785, all in and out behind: 1.34M out 8 / 35
2,000,000 behind: 1.91M in, 105k out behind: 1.50M in, 1.23M out 9 / 59

Queen’s own clients are the 2.0.0-beta.2 runs of throughput and latency, driven from four load machines. On the facade build with the same three load machines as the Kafka clients, they measured 89 / 212 ms at a million (recorded in b51e4419’s message): our HTTP load generator needs more CPU per message than franz-go, and on three machines it adds latency of its own.

Partition count

At 1,000,000 msg/s, p50 / p99 in milliseconds.

Partitions Kafka clients through Queen Queen’s own clients Kafka 4.3.1
200 44 / 108 61 / 146 8 / 18
10,000 45 / 142 51 / 144 15 / 46
50,000 78 / 200 49 / 103 692 / 1,466
100,000 134 / 251 50 / 163 behind: 2.7 s / 10.9 s
p99 end-to-end latency at 1,000,000 msg/s against the partition count, log scales. Kafka clients through Queen's Kafka facade: 108 ms at 200 partitions, 142 at 10,000, 200 at 50,000 and 251 at 100,000. Queen's own clients: between 103 and 163 ms. Kafka 4.3.1 itself: 18 ms at 200 partitions, 46 at 10,000, 1.5 s at 50,000, and behind at 100,000 with a p99 of 10.9 s.p99 end-to-end latency10 ms100 ms1 s10 s1k10k100kpartitionsKafka clients on QueenQueen clientsKafka 4.3.1
The same Kafka clients, pointed at Queen instead of Kafka. Under 10,000 partitions Kafka is faster; past 50,000, Queen's partition-per-entity store keeps the Kafka clients under a quarter of a second while Kafka falls behind. Source: benchmark-queen/2026-09-30-stage03-3node, benchmark-queen/2026-09-30-kafka-pulsar

With many topics of 10,000 partitions in total, Kafka clients through Queen measured 44 / 140 ms on 10 topics of 1,000, 58 / 165 ms on 100 of 100, and 41 / 100 ms on 1,000 of 10.

CPU

Cores used by Queen’s processes on the three brokers at 1M msg/s:

Partitions Kafka clients through Queen Queen’s own clients
200 22.1 17.2
10,000 25.2 18.4
50,000 36.9 18.0
100,000 49.7 18.3
Cores used by Queen at 1,000,000 msg/s, summed over three brokers, with Kafka clients through the Kafka facade and with Queen's own clients. 200 partitions: 22.1 and 17.2. 10,000: 25.2 and 18.4. 50,000: 36.9 and 18.0. 100,000: 49.7 and 18.3.Kafka clients on QueenQueen clients02040cores, three brokers200 partitions10,00050,000100,00022.117.225.218.436.91849.718.3
The gap grows on the fetch path. The facade answers Fetch below fetch sessions, so every fetch names all of its consumer's partitions; a Queen pop asks for work and is handed partitions that have some. Source: benchmark-queen/2026-09-30-stage03-3node

The gap that grows with partitions is on the fetch path. The facade answers Fetch up to version 6, below fetch sessions (KIP-227), so every fetch names all of its consumer’s partitions, and at 100,000 partitions the nine client processes used 22 cores against Queen and 6 against Kafka. We left sessions out on purpose: a session is per-connection state, and the facade keeps none, so it can restart the way a Kafka broker does.

What made it fast

On 2026-09-30 the same load at 1M msg/s on 200 partitions got through the facade at 935,000 in and 924,000 out, with an e2e of 2.0 s at p50 and 3.8 s at p99 (run qk-1x200-1000k). Four changes, all in b51e4419, took that to 44 / 108 ms:

  • Produce and Fetch call the broker’s push and fetch as typed calls: the same commands, admission and reads as POST /api/v1/push and POST /api/v1/fetch, with no JSON in between. A test pins that both paths give byte-equal answers, so the two cannot drift apart unnoticed.
  • A Kafka record travels as bytes in both directions.
  • A topic created through Kafka gets dedupWindowSeconds: 0. A Kafka record carries no Queen transactionId, so the window could only ever check a brand-new id; Kafka’s own duplicate suppression, the idempotent producer’s sequence window, does not need it.
  • Produce requests that a connection already has buffered are written as one push and answered in order. QUEEN_KAFKA_COALESCE (default 5, Kafka’s max.in.flight.requests.per.connection) caps how many, and nothing waits for a request that has not arrived.

On the final build, Kafka Streams (WordCountDemo) and Kafka Connect in distributed mode, with a file source into a topic and a file sink out of it, passed four runs out of four on fresh three-node clusters.

Still open

Overload is the first open problem. At 2M msg/s offered the facade accepted up to 1.9M a second and its consumers starved: 105,000 consumed in the archived run, and anywhere from 42,000 to 1,550,000 in repeats with and without produce coalescing. Queen’s own clients at the same rate shed at their in-flight cap and kept consuming 1.23M. We do not know the mechanism yet.

Long runs are the second. A 10-minute run at 1M msg/s (qk3, 2026-10-01) moved 595 million messages in and out with no errors, and its last 30 seconds held 51 / 111 ms. Three stretches of 20 to 100 seconds, around minutes 3, 6 and 8, had produce p99 of 2 to 8 seconds and put the run’s e2e p99 at 10 to 12 seconds. The first began when free memory on the nodes ran out and the followers fell into direct reclaim, while each node’s anonymous memory grew by about 13 GB over the ten minutes.

The facade is also outside the Jepsen tests, which drive the native HTTP API; it has a compatibility suite of its own.

How to point a Kafka client at Queen, and what the facade supports, is in the Kafka guide.

Navigation

Type to search…

↑↓ navigate↵ selectEsc close