Skip to content

How we measure

The machines, the load generators and what they count, each system's settings and durability, Queen's non-default settings, the commands that ran every benchmark in this section, and the archive the runs live in.

Updated View as Markdown

Every number in this section comes from a folder of logs, and most can be produced again with one command. This page is how: the machines, the load generators and what they count, the settings each system ran with, and the scripts.

The machines

Role Count Shape
Brokers 3 DigitalOcean, 16 vCPU, 31 GB, one local disk with ext4. Two Xeon Gold 6548N, one Xeon Platinum 8358
Load machines 3, and a fourth for Queen’s native runs 16 vCPU, 31 GB, Xeon Platinum 8358

All of them ran Ubuntu 24.04 in one VPC, with the same kernel settings on every machine for every system (common/os-tune.sh): swappiness 1, larger socket buffers and backlogs, a million open files, noatime, transparent huge pages left at Ubuntu’s madvise, and no IRQ pinning. Clocks were synced with chrony, each within a third of a millisecond of NTP. Only one system ran on the brokers at a time, and every start script refuses to start while another system’s ports are open.

The older Platinum broker caps a raft leader about 15 to 20% lower, so Queen’s harness stops that node whenever it leads and starts it again as a follower, until another node leads. Kafka, Redpanda and Pulsar spread partition leadership over all three brokers, the older one included.

The load generators

Queen’s native runs used goload, a Go program built on Queen’s Go client. Kafka, Redpanda and Pulsar used kload (franz-go) and pload (pulsar-client-go), and the transactional runs added qload for Queen. The three share one core copied from goload, so all four systems are driven and measured the same way.

The producers run open loop. Each unit, 100 messages for one partition, leaves at its scheduled instant, and a unit that would pass the in-flight cap of 5,000 per process is shed and counted, never queued, so a broker that slows down shows up as latency and shedding instead of as a smaller load. Producers never retry, since a retry would offer the same messages twice. Produce latency runs from a unit’s scheduled instant to its last acknowledgement, so time spent waiting to be sent counts.

Every message carries the instant its unit was scheduled, and e2e is the consumer’s receive time minus that instant, taken before the ack. The two ends are on different machines, so e2e includes the clock offset between them, a fraction of a millisecond here; the Kafka and Pulsar loaders also print a same-machine e2e with no offset at all, and it agreed with the cross-machine one. All of them send the same JSON payload, about 256 bytes (289 on average), and keep latencies in log-linear histograms with the same buckets, so their percentiles compare.

How one unit of load is timed. A producer schedules a unit of 100 messages for one partition at instant t0. It is sent if the process has fewer than 5,000 units in flight, otherwise it is shed and counted, never queued. The broker acknowledges it, and produce latency is the acknowledgement time minus t0. A consumer later receives the messages, which carry t0, and e2e latency is the receive time minus t0, taken before the ack.producerload machinebrokerthree nodesconsumerload machineunit scheduled at t0:100 messages, one partitionpush, if under 5,000 in flightotherwise shed and countedackproduce latency = ack − t0pop or fetchmessages carrying t0e2e = receive − t0,before the ack
Both latencies start at the scheduled instant, not at the send, so a broker that slows down shows up as latency and shedding instead of as a smaller load. Source: benchmark-queen/2026-09-30-kafka-pulsar (goload, kload, pload)

Every 10 seconds each process prints one line. A figure “over the last 30 seconds” is the three full windows before the process’s final one: rates are summed over the processes, a latency is the worst process’s. CPU is the system’s own processes, from 2-second samples over the loaded window, summed over the three brokers. A system fell behind when it consumed under 97% of the offered rate or its e2e p99 passed 5 seconds.

Queen’s settings

Each benchmark binary was built with Cargo’s fastrel profile, which is release without link-time optimization and with 16 codegen units, so it builds in minutes between runs; the published images use release, with thin LTO and one codegen unit. The cluster had three voters on the openraft replicator, with the data directory on the local disk and authentication and tenancy headers off. Every run used one tenant, so one raft group. Four settings were set explicitly, and three of them differ from the defaults (all are in the configuration reference):

Setting Default These runs Why
QUEEN_RAFT_LOG_CACHE_MB 512 4096 Applied entries leave the leader’s memory past this size even if a follower still needs them, and a follower behind the cache is fed from disk. At a million messages a second, 512 MB is a second or two of log.
QUEEN_RAFT_TXN_WINDOW_MIN_S 900 300 The least time a message’s hash list and its queue-log record are kept, so no log file is freed sooner. The leader wrote about 250 MB/s at a million messages a second, so 900 seconds of log would pass 200 GB a node.
QUEEN_QLOG_SHARDS 0 4 Every queue’s records go to one of four shared logs, so a group commit fsyncs at most four files however many queues it touches.
QUEEN_LANES 8 8; 16 in the transaction runs The planning threads that plan pushes in parallel.

The throughput and partition runs had deduplication off (dedupWindowSeconds: 0) and the soak a 60-second window. Completed messages were kept 1,800 seconds in the 60-second runs, so retention deleted nothing during a point, and 300 seconds in the soak. Leases were 30 seconds. Consumers popped up to 10 partitions and 1,000 messages at a time, long-polling for 2 seconds, and acked asynchronously with up to 256 acks in flight per process; in the soak the leader chose the pop width. There were 300 consumers on one queue, 1,008 at 1,000 queues and 10,008 at 10,000, where the other systems ran 198 consumers on 200 partitions and 297 above. Partitions were created before each timed run by pushing one message into each and draining it.

The other systems’ settings

SPEC.md in the Kafka and Pulsar folder is the contract the runs followed, with every setting and the reason for it, and the rule that tuning may change performance knobs only, never durability or ordering. In short:

  • Kafka 4.3.1, KRaft with three combined broker and controller nodes. Replication factor 3, min.insync.replicas 2, acks=all, idempotent producers, lz4, 5 ms linger, 8 replica fetchers, 8 network and 16 I/O threads, message.max.bytes 16 MiB, a 6 GB heap (10 GB from 50,000 partitions) on G1, KIP-848 consumer groups, no fsync (page cache, Kafka’s production model).
  • Redpanda 26.2.3, production mode with rpk redpanda tune all and rpk iotune, write caching off, classic consumer groups (26.2.3 does not answer KIP-848), its partition memory limits raised for 75,000 partitions and more. Its tuners were undone before any other system ran.
  • Pulsar 4.2.4, ZooKeeper, a bookie and a broker on each node, ensemble 3, write quorum 3, ack quorum 2, journal fsync on, 48 namespace bundles, generational ZGC, producer batching of 5 ms or 1,000 messages with lz4, failover subscriptions.

What each system acknowledges after, side by side, is on throughput and latency.

Running it again

Each folder below has the scripts that ran and a README or SPEC.md. Fill in the machines’ addresses first: cluster.env for Queen’s harness, hosts.env for the others.

Queen’s matrix and single points, from benchmark-queen/2026-09-30-stage03-3node/, with the broker binary in /root/qr/ on each broker:

./setup30.sh                    # scripts and goload onto every machine
./campM.sh queen-b2 mb2c        # the 1M matrix: 60 s per point, deduplication off, pops of 10 partitions
DEDUP=0 GX="-pop-partitions 10" NL=4 RATE=1000000 ./run30.sh mb2c-q1-t10m 1 10000000 queen-b2 60
python3 grid_report.py runs mb2c-q1-t200-r1m mb2c-q1-t10m

The soak, from the same folder:

BIN=/root/qr/queen-b1 SECS=36000 ./soak.sh soak6   # 10 hours, a fault every 90 minutes
./soak-status.sh soak6                             # the latest window of every process
./soak-collect.sh soak6                            # fetch everything into runs/soak6

Kafka, Redpanda, Pulsar and the transactions, from benchmark-queen/2026-09-30-kafka-pulsar/:

(cd mqload && GOWORK=off ./build.sh)
common/deploy.sh harness tune kafka pulsar redpanda clock check
kafka/grid.sh && report/report.py runs/kafka/* --loaders
redpanda/grid.sh && redpanda/cluster.sh stop && redpanda/cluster.sh untune
pulsar/grid.sh
txn/deploy.sh && txn/grid.sh && report/report.py --txn runs/txn/*/*

The archive

Folder under benchmark-queen/ What it holds
2026-09-30-stage03-3node/ Queen’s harness for 2026-09-30 and 10-01 and every Queen run on these pages: the 2.0.0-beta.2 matrix (runs/mb2c-*) and the soaks (runs/soak1 to soak5), with the goload source
2026-09-30-kafka-pulsar/ SPEC.md, the runbook, the mqload loaders, each system’s install and configuration scripts, report/report.py, and the Kafka, Redpanda, Pulsar, transaction and Kafka-client runs
2026-09-29-1m-3node/ the 2026-09-29 runs at 1M msg/s on 500,000 partitions, deduplication on and off, and their harness
2026-09-27-jepsen-archive/ Jepsen P8 to P10: every run, the Jepsen stores and the binaries they tested

Three older 2.0 folders are kept and not quoted, because the engine they measured has been rewritten since: 2026-09-23-raft-32core/ (one 32-vCPU node), 2026-09-24-raft-3node/ (Queen and Kafka by partition count at 300,000 msg/s on 32-vCPU nodes) and 2026-09-25-queue-count/ (one node, the queue-count problem that led to shared queue logs). Results of Queen 1.x, whose storage was PostgreSQL, are in the git history of this site.

Navigation

Type to search…

↑↓ navigate↵ selectEsc close