Skip to content

Partition count

One queue from 200 to 10,000,000 partitions at 1,000,000 msg/s on three nodes, with Kafka, Redpanda and Pulsar on the same machines up to 100,000: latency, CPU, what a partition costs Queen in memory and creation time, and what many queues cost.

Updated View as Markdown

This is the test Queen was built for. Give every entity its own partition and your partition count is your entity count: a million customers, ten million orders. We offered one queue 1,000,000 msg/s at partition counts from 200 to 10,000,000, and ran Kafka, Redpanda and Pulsar on the same machines for as far as each would go.

The shape

The load is the one on throughput and latency: 256-byte JSON messages pushed 100 at a time to one partition, 1,000,000 msg/s offered, 60 seconds per point, figures over the last 30. Every partition existed before the timed run. Queen’s load generator pushed one message into each partition and drained it; Kafka and Redpanda topics of 2,000 partitions and more were warmed with one record per partition; Pulsar’s subscription was created on every partition before any producer sent. No system was measured while it created partitions. Queen 2.0.0-beta.2 ran on 2026-10-01 with deduplication off (tags mb2c-q1-t* in benchmark-queen/2026-09-30-stage03-3node/runs/); the others are in benchmark-queen/2026-09-30-kafka-pulsar/runs/.

Latency at 1,000,000 msg/s

End-to-end latency, p50 / p99 in milliseconds, the worst of the load processes. “Behind” marks a point where the system consumed under 97% of the offered rate or its p99 passed 5 seconds.

Partitions Queen Kafka Redpanda Pulsar
200 61 / 146 8 / 18 16 / 187 7 / 21
2,000 not run 9 / 194 18 / 105 7 / 15
10,000 51 / 144 15 / 46 371 / 1,810 7 / 32
50,000 49 / 103 692 / 1,466 behind: 881k out, p99 21 s behind: 232k out, p99 47 s
100,000 50 / 163 behind: p50 2.7 s, p99 10.9 s behind: 236k out, p99 47 s behind: 199k out, p99 47 s
1,000,000 46 / 123 not run not run setup failed after 31 min
5,000,000 45 / 138 not run not run not run
10,000,000 51 / 122 not run not run not run
p99 end-to-end latency against partition count, both on log scales, at 1,000,000 msg/s offered on three nodes. Queen stays between 103 and 163 ms from 200 to 10,000,000 partitions. Kafka goes from 18 ms at 200 partitions to 1.5 s at 50,000 and falls behind at 100,000. Redpanda reaches 1.8 s at 10,000 and falls behind from 50,000. Pulsar stays near 20 ms up to 10,000 and falls behind from 50,000.p99 end-to-end latency10 ms100 ms1 s10 s1k10k100k1M10Mpartitions in the queueQueenKafkaRedpandaPulsar
p99 end-to-end latency with 1,000,000 msg/s offered to one queue. A hollow mark is a point where the system fell behind: it consumed under 97% of the offered rate or its p99 passed 5 seconds. Kafka, Redpanda and Pulsar were not run past 100,000 partitions. Source: benchmark-queen/2026-09-30-stage03-3node, benchmark-queen/2026-09-30-kafka-pulsar

Queen consumed the full million a second at every count. Three points need a footnote. At 50,000 partitions, three of Redpanda’s nine load processes were killed partway through (exit status 137), so that point is partly the load generator’s. Pulsar’s points at 50,000 and 100,000 ran with client-side batching off, one message per batch, because the Go client cannot keep that many batching producers in one load process, which makes them harder on Pulsar than its other points; six of the nine load processes at 50,000 and three at 100,000 crashed with a nil-pointer panic inside the Pulsar Go client (v0.21.0), and the brokers were at 38 of their 48 cores. Queen at 10,000,000 had one slow 10-second window early in the run, which puts the p99 over the whole run at 440 ms; the last 30 seconds held at 122 ms.

CPU

Cores used by each system’s own processes, summed over the three brokers.

Partitions Queen Kafka Redpanda Pulsar
200 17.2 9.7 14.2 7.1
2,000 not run 15.3 16.2 7.8
10,000 18.4 28.3 32.5 10.1
50,000 18.0 29.2 35.1, fell behind 37.4, fell behind
100,000 18.3 32.4, fell behind 38.7, fell behind 38.2, fell behind
1,000,000 18.3
10,000,000 17.3
Cores used at 1,000,000 msg/s offered against the partition count, log scale, summed over three brokers. Queen stays between 17.2 and 18.4 cores from 200 to 10,000,000 partitions. Kafka goes from 9.7 cores at 200 partitions to 28.3 at 10,000 and 32.4 at 100,000, where it had fallen behind. Redpanda goes from 14.2 to 32.5 at 10,000 and 38.7 at 100,000. Pulsar goes from 7.1 to 10.1 at 10,000 and 38.2 at 100,000.cores, three brokers0102030401k10k100k1M10Mpartitions in the queueQueenKafkaRedpandaPulsar
Queen's cost does not move with the partition count: a partition is a few rows in memory. The others keep machinery per partition, which is cheap at 200 and grows with every one added. Hollow marks: the system had fallen behind. Source: benchmark-queen/2026-09-30-stage03-3node, benchmark-queen/2026-09-30-kafka-pulsar

Redpanda’s column is its count of busy shard time; the operating system counts more, because its reactors poll.

Why the line stays flat

In Queen a partition is a few rows in the store, which every node keeps in memory: the partition itself, its offsets, and a cursor for each consumer group that reads it. Its messages are records in its queue’s log, beside every other partition’s, and one fsync covers everything a group of entries wrote to that log. So a new partition adds memory and nothing else, neither a file nor a replication stream nor an fsync of its own, and the leader’s CPU at ten million partitions is what it was at two hundred.

Kafka keeps a log of its own for every partition on every replica, with its own files and its own place in replication; Redpanda runs a raft group per partition; Pulsar keeps a managed ledger per partition in BookKeeper. That machinery is cheap with a few hundred partitions and grows with every one you add, which is what the CPU table shows from 10,000 on.

What a partition costs Queen

Partitions Leader memory (anonymous RSS) Time for the pushes that created them
200 3.6 GB under a second
100,000 3.9 GB 0.6 s
1,000,000 5.6 GB 6.2 s
5,000,000 11.8 GB 28.6 s
10,000,000 19.9 GB 60.8 s
The Queen leader's anonymous memory against the number of partitions in the queue, 3.6 GB with 200 partitions, 3.9 GB with 100,000, 5.6 GB with 1,000,000, 11.8 GB with 5,000,000 and 19.9 GB with 10,000,000.leader memory (anonymous RSS)0 GB5 GB10 GB15 GB20 GB02.5M5M7.5M10Mpartitions in the queueQueen leader
A straight line: memory is what a partition costs, about 2 KB each on the leader, and every node holds the same store, so memory sets the partition ceiling. Source: benchmark-queen/2026-09-30-stage03-3node

That is about 2 KB per partition on the leader, and every node holds the same store, so memory is what sets the partition ceiling: ten million partitions took two thirds of a 31 GB node. The first push that names a partition creates it, in the same log entry as its message, so creating them ran at about 165,000 partitions a second with nothing declared beforehand. For comparison, Kafka took 255 seconds to create a 100,000-partition topic (until every replica’s log was on disk) and another minute to warm it, and Pulsar’s setup for 1,000,000 partitions failed after 31 minutes.

Many queues

A queue costs more than a partition. One million partitions spread over many queues, at the same 1M msg/s:

Queues × partitions In / out e2e p50 / p99
1 × 1,000,000 1.00M / 1.00M 46 / 123 ms
1,000 × 1,000 1.00M / 1.00M 95 / 216 ms
10,000 × 100 1.00M / 1.01M 185 / 561 ms; 76 pops timed out at their 2 s deadline
End-to-end latency at 1,000,000 msg/s with one million partitions arranged three ways. One queue of 1,000,000 partitions: p50 46 ms, p99 123 ms. 1,000 queues of 1,000: p50 95 ms, p99 216 ms. 10,000 queues of 100: p50 185 ms, p99 561 ms, with 76 pops timing out at their 2-second deadline.p99p500 ms200 ms400 msend-to-end latency, ms1 × 1,000,0001,000 × 1,00010,000 × 100123 ms46 ms216 ms95 ms561 ms185 ms
A queue costs more than a partition. With about one consumer per queue, each pop carried fewer messages at 10,000 queues (about 205 against 450) and half came back empty. Source: benchmark-queen/2026-09-30-stage03-3node

The run at 10,000 queues was repeated the same day at 167 / 518 ms. There were 10,008 consumers for 10,000 queues, about one each, so a pop carried about 205 messages against about 450 with one queue, half the pops came back empty, and an ack took 60 ms against 20. Kafka and Pulsar ran their many-topic shapes at 10,000 and 100,000 partitions in total, so the totals differ by ten times or more: with 1,000 topics of 10 partitions Kafka held 20 / 79 ms and Pulsar 8 / 34 ms, and with 1,000 topics of 100 partitions Kafka fell behind (p99 5.8 s).

Keys hashed onto partitions

If all you need is order per key, Kafka and Pulsar can hash the keys onto a fixed set of partitions, and they are very good at it. Kafka spread 10,000,000 keys over 1,000 partitions at 8 / 13 ms, Pulsar 10,000,000 keys over 48 partitions with a Key_Shared subscription at 12 / 21 ms, and Redpanda, laid out as Kafka, at 64 / 293 ms. What a hash gives up is isolation: a slow key holds up every key that shares its partition, and you cannot give a consumer one key or ask where one key stands. Queen’s partition per entity keeps that, and the tables above are what it costs.

What this page does not show

Traffic on every partition at once. Each push carried 100 messages for one partition, so 10,000 pushes a second reached at most 600,000 partitions in a minute: all ten million were created, written and drained in the prefill, and the timed minute touched 600,000 of them. Traffic on all of them within a minute means pushes of a message or two per partition, the shape described on throughput and latency.

Each point also ran only 60 seconds (the long runs are on soak and failover), and we did not run Queen at 2,000 partitions or the other systems above 100,000.

Navigation

Type to search…

↑↓ navigate↵ selectEsc close