Your Kafka producers, consumers and admin tools can use Queen as they are: change
bootstrap.servers and nothing else. A topic is a Queen queue and a consumer group is a Queen
consumer group, so Kafka clients and Queen clients read and write the same streams, and you can move
an application over one service at a time.
Try it
docker run -d --name queen --platform linux/amd64 -p 6632:6632 -p 9092:9092 \
-v queen-data:/var/lib/queen/raft \
-e QUEEN_KAFKA_EMBEDDED=true \
-e QUEEN_KAFKA_ADVERTISED_ADDR=localhost:9092 \
ghcr.io/queen-mq/queen:latestQUEEN_KAFKA_EMBEDDED=true starts the Kafka listener inside the broker, on port 9092.
QUEEN_KAFKA_ADVERTISED_ADDR is the address Kafka clients come back to after their first request,
so it has to be one they can reach. It has no default, and a node started without it refuses to
boot and tells you why.
Send three keyed records with kcat and read them back:
printf 'customer-123:{"orderId":8891,"total":42}\ncustomer-456:{"orderId":8892,"total":17}\ncustomer-123:{"orderId":8893,"total":5}\n' \
| kcat -P -b localhost:9092 -t orders -K:
kcat -C -b localhost:9092 -t orders -o beginning -e -q -f '%t [%p] offset %o key=%k %s\n'orders [307] offset 0 key=customer-123 {"orderId":8891,"total":42}
orders [307] offset 1 key=customer-123 {"orderId":8893,"total":5}
orders [400] offset 0 key=customer-456 {"orderId":8892,"total":17}Nobody created orders. kcat’s first metadata request did: the queue came into existence 1,024
partitions wide, kcat hashed each key onto one of them, and both orders of customer-123 landed in
partition 307 in the order they were sent. Queen’s HTTP API reads the same records:
curl -s 'localhost:6632/api/v1/pop/queue/orders?batch=10' | jq -c '.partition, .messages[].data'"307"
{"k":"Y3VzdG9tZXItMTIz","t":1790948151426,"v":"eyJvcmRlcklkIjo4ODkxLCJ0b3RhbCI6NDJ9"}
{"k":"Y3VzdG9tZXItMTIz","t":1790948151427,"v":"eyJvcmRlcklkIjo4ODkzLCJ0b3RhbCI6NX0="}The partition is the Kafka partition and the offsets are the Kafka offsets. Each record is one
message. Its payload carries the key (k) and the value (v) in base64, because a Kafka key or
value is bytes and a Queen payload is JSON, then the producer’s timestamp (t) and, when there are
any, the headers (h).
Consumer groups are shared the same way. Run a group consumer, and stop it with Ctrl-C once it has printed the three records:
kcat -b localhost:9092 -G billing -o beginning -q -f 'billing [%p] offset %o key=%k\n' ordersbilling [400] offset 0 key=customer-456
billing [307] offset 0 key=customer-123
billing [307] offset 1 key=customer-123curl -s localhost:6632/api/v1/consumer-groups \
| jq -c '.[] | select(.name == "billing") | {name, queueName, partitionCursors, totalLag}'{"name":"billing","queueName":"orders","partitionCursors":2,"totalLag":0}billing is a Queen consumer group now. Its commits are the group’s positions, its lag is in the
dashboard beside every other group, and a Queen consumer that joins it starts where the Kafka
consumer stopped. Push one more order over HTTP and pop it as billing:
curl -s -X POST localhost:6632/api/v1/push -H 'content-type: application/json' \
-d '{"items":[{"queue":"orders","partition":"307","payload":{"orderId":8894,"total":99}}]}' > /dev/null
curl -s 'localhost:6632/api/v1/pop/queue/orders/partition/307?consumerGroup=billing&batch=10' \
| jq -c '.messages[] | {offset, data}'{"offset":2,"data":{"orderId":8894,"total":99}}And a Kafka consumer reads what a Queen client pushed as a record with no key, whose value is the JSON document:
kcat -C -b localhost:9092 -t orders -p 307 -o 2 -c 1 -q -f 'offset %o key=%k value=%s\n'offset 2 key= value={"orderId":8894,"total":99}How it works
The facade runs inside the broker process on threads of its own (QUEEN_KAFKA_THREADS,
min(4, cores / 2) by default) and keeps no storage of its own. A produce request becomes one push
(produce requests that arrive together on a connection can share one), and a fetch request one read
of up to 1,024 partitions, more in chunks of 1,024 one after another. Both go through the same
admission, the same raft log and the same reads as POST /api/v1/push and POST /api/v1/fetch,
called in memory with no HTTP request and no JSON body. The rarer calls, commits, metadata and time
lookups, go through the same router with a JSON body.
We kept Kafka on Queen’s one pipeline on purpose: a record is a Queen message from the moment it is
written, so replication, retention, the dashboard and native consumers all apply to it, and there is
no second storage path to keep in step with the first. (We built one that stored Kafka batches
verbatim, and removed it in September 2026 because every change to the main pipeline broke it.)
server/src/kafka_inproc.rs, protocols/queen-kafka/src/handlers/produce.rsWith the facade on, every node of a cluster runs one, and every node serves produce and fetch for every partition. The raft leader assigns the offsets for the whole cluster, so two nodes appending to one partition can never hand out the same offset.
acks=1 and acks=all take the same path: a produce is answered when its push is committed to the
raft log, which on a cluster means a majority of the voters have it. acks=0 writes the same way
and sends nothing back. Kafka clients keep several produce requests in flight on a connection (five
by default in Java), and the facade writes the ones already waiting in the socket as one push, then
answers each in order (QUEEN_KAFKA_COALESCE, default 5). One push also keeps them in order on
every node, which is what an idempotent producer counts on.
Topics created through Kafka get no dedup window (dedupWindowSeconds: 0). A Kafka record carries
no Queen transactionId, so the broker gives each one a fresh id, and a dedup check could only ever
ask whether a brand-new id had been seen before. Kafka’s own duplicate suppression, the idempotent
producer, is below. A queue that already existed keeps its own settings.
The flip side: if Queen clients also push into a topic Kafka created, with their own
transactionIds, repeats are stored twice until you give the queue a window
(dedupWindowSeconds, see queue options).
How Kafka maps onto Queen
| Kafka | Queen |
|---|---|
| Topic | A queue of the same name, created by CreateTopics or by the first metadata request that names it |
Partition n |
The partition named "n" |
| Partition count | The topic’s width: what CreateTopics declared, else QUEEN_KAFKA_DEFAULT_PARTITIONS (1,024), or more if more partitions exist |
| Record | One message: key, value, headers and timestamp in a JSON envelope |
| Offset | The partition’s offset, assigned by the raft leader |
| Consumer group | The Queen consumer group of the same name |
| Committed offset | The group’s position on that partition |
acks=1, acks=all |
One path: answered once the push is committed to the raft log |
| Replication factor | Accepted and reported as 1; every raft voter holds every partition |
retention.ms |
The queue’s retention, in whole seconds |
cleanup.policy=compact |
Retention off: every record is kept |
| Idempotent producer | A sequence window per producer and partition, in the node’s memory |
| Transaction | Staged in the node’s memory, committed with its offsets as one raft entry |
Partitions and width
A Kafka client spreads keys by hashing them modulo the partition count, so a topic needs a fixed
width. A Queen queue has none: a partition is a name that exists from its first message. The facade
keeps a width per topic and advertises the larger of that width and the partitions that exist. The
width is what CreateTopics asked for, so kafka-topics.sh --create --partitions 8 makes an
eight-partition topic, and QUEEN_KAFKA_DEFAULT_PARTITIONS for a topic created on first use. A
topic never narrows, because a smaller count would send existing keys to other partitions, and
CreatePartitions widens a topic created through Kafka. CreateTopics writes the width before it
creates the queue, so on a cluster every node reports the declared width from its first answer
about a new topic. Kafka Connect, which reads the partitions of the topics it has just created
straight back from metadata, depends on exactly that.
The default is wide because partitions are cheap here: until something is written to it, a
partition costs the broker a key and a counter, and a wide topic leaves a consumer group room to
grow. Width has one real cost, the metadata answer, which lists every partition at about 50 bytes
each, so one topic is capped at QUEEN_KAFKA_MAX_PARTITIONS (1,000,000).
Queen partitions whose names are not numbers, such as customer-123, count toward the width, but a
Kafka client cannot address them. A produce through Kafka always creates numbered partitions, so
this matters only for queues that Kafka and Queen clients share.
Consumer groups and committed offsets
Groups use Kafka’s classic protocol (JoinGroup, SyncGroup and heartbeats), and the assignment is
computed by the client, so range, round-robin, sticky, cooperative-sticky and the Kafka Streams
assignor all work. The first member of a new group waits three seconds for others before the first
assignment (QUEEN_KAFKA_GROUP_JOIN_DELAY_MS, Kafka’s group.initial.rebalance.delay.ms).
A committed offset is the group’s position on the partition: one raft entry per commit of up to
4,096 partitions, more split into several (and a group’s first commit on a node also writes a small
index row). That is why billing showed up in /api/v1/consumer-groups with its lag,
and why a Queen consumer of billing resumed at offset 2. Kafka consumers read with fetch requests
and take no leases, so Queen’s redelivery, retry budget and dead-letter queue do not apply to them;
a Kafka consumer handles its own failures, as it does on Kafka. Both kinds of consumer share a
group’s position, but not its work. A Kafka commit sets the position the way a seek does and
releases any lease a Queen consumer holds on that partition, so hand a group from one kind of
consumer to the other, and do not run both at once.
Committed offsets never expire. Kafka drops the offsets of a group that has been empty for seven
days (offsets.retention.minutes), and a consumer that comes back after that starts over from
auto.offset.reset without a word; here a group keeps its position until you delete it.
kafka-consumer-groups.sh --delete and --delete-offsets do that, and since a Kafka group is a
Queen group, deleting one from a Kafka tool deletes the Queen consumer group of that name, on every
queue.
Idempotent producers
Java clients since 3.0 and franz-go turn idempotence on by default, so a producer with no settings at all is an idempotent one, and the facade keeps Kafka’s contract with it. A batch sent again after a lost answer is answered with the offsets it got the first time and is not written twice. A batch that would leave a gap in the sequence is refused, and nothing is written.
The sequence windows live in the memory of the node a producer writes to, one per producer and
partition, up to QUEEN_KAFKA_MAX_PRODUCER_STATES (4,194,304 of them, about 210 bytes each). Kafka
keeps producer state in its log. Doing that here would put a dedup key on every record of every
idempotent producer, a cost on the write path we have not taken, so after a node restart, an
eviction or a connection that moves to another node, the producer’s next batch is refused as out of
sequence. The client recovers by bumping its epoch (KIP-360), and at most the five batches it had in
flight can be written twice.
Transactions
A transactional producer needs no extra settings, and sendOffsetsToTransaction in a
consume-transform-produce loop works too. Against one node, a transaction of three records behaves
like this: each produce inside it is answered with offset -1, the topic’s end offset does not move
while the transaction is open, the commit moves it by exactly three, and an aborted transaction
leaves it where it was.
The records wait in the memory of the node that holds the transaction. commitTransaction() writes
them, together with the consumer offsets the transaction carries, as one raft entry, so the records
and the offsets land together or not at all, and abortTransaction() throws the stage away. No
uncommitted record ever enters the log, which is why read_committed costs nothing here and why a
committed transaction of N records advances the log by N, where Kafka adds a commit marker. Two
producers with the same transactional.id fence each other through an epoch kept in Queen’s
key-value store, and the fenced producer’s commit writes nothing.
Because a transaction lives in memory until it commits, it has limits Kafka does not: 8 MiB and
50,000 records per transaction, 128 MiB and 1,024 open transactions per node, at most 200 partitions
and 62 offsets each, and a timeout of at most an hour (the QUEEN_KAFKA_TXN_MAX_* settings in the
reference). A stage outlives the connection that opened it until its own
timeout, so a commit resumed from a new connection still lands; that is what Flink’s exactly-once
sink does after a task manager fails over. A restart of the node holding the stage is a different
matter: the resumed commit is refused with INVALID_TXN_STATE, and the transaction’s records are
never written.
Compaction and retention
retention.ms on a topic created through Kafka becomes the queue’s retention, rounded down to whole
seconds; -1, the default for these topics, keeps records forever. cleanup.policy=compact is
accepted, and Queen keeps every record of such a topic. Retention is off whatever retention.ms
says, so the last value of every key is always in the log, next to every earlier one. A reader that
replays a compacted topic, a Kafka Streams changelog or Kafka Connect’s config topic, rebuilds the
same state from the whole log as from a compacted one. It reads more, and the topic grows until
real compaction exists; tombstones are kept like any other record. compact,delete keeps the
retention and compacts nothing. What Kafka Streams, Connect and Flink
ask for is on a page of its own.
Run it on a cluster
Each node runs its own facade, and left alone each one presents itself as a one-broker Kafka
cluster: produce and fetch work against any node, but the members of one consumer group must
all connect to the same node, because each node coordinates its own groups. Cluster mode turns the
nodes into one Kafka cluster. Give every node a distinct QUEEN_KAFKA_NODE_ID from 1 to 64, its own
reachable QUEEN_KAFKA_ADVERTISED_ADDR, and a QUEEN_TOKEN, which the facade writes its node
registry with (with broker authentication off, any value will do). For the three-node Compose file
in deploy/compose/three-node/ (Run a cluster), this override sits next to
compose.yaml in the same folder:
# docker compose -f compose.yaml -f compose.kafka.yaml up -d
x-kafka: &kafka
QUEEN_KAFKA_EMBEDDED: "true"
QUEEN_TOKEN: kafka-facade
services:
queen-1:
environment:
<<: *kafka
QUEEN_KAFKA_NODE_ID: "1"
QUEEN_KAFKA_ADVERTISED_ADDR: 127.0.0.1:16092
ports: ["127.0.0.1:16092:9092"]
queen-2:
environment:
<<: *kafka
QUEEN_KAFKA_NODE_ID: "2"
QUEEN_KAFKA_ADVERTISED_ADDR: 127.0.0.1:26092
ports: ["127.0.0.1:26092:9092"]
queen-3:
environment:
<<: *kafka
QUEEN_KAFKA_NODE_ID: "3"
QUEEN_KAFKA_ADVERTISED_ADDR: 127.0.0.1:36092
ports: ["127.0.0.1:36092:9092"]The addresses say 127.0.0.1 because the ports are published on that address only, and some
systems (macOS among them) try localhost as ::1 first. Produce through node 1 and ask node 2
what the cluster looks like:
docker compose -f compose.yaml -f compose.kafka.yaml up -d
echo 'customer-123:{"orderId":8891}' | kcat -P -b 127.0.0.1:16092 -t orders -K:
kcat -L -b 127.0.0.1:26092 -t orders | head -9Metadata for orders (from broker -1: 127.0.0.1:26092/bootstrap):
3 brokers:
broker 1 at 127.0.0.1:16092 (controller)
broker 2 at 127.0.0.1:26092
broker 3 at 127.0.0.1:36092
1 topics:
topic "orders" with 1024 partitions:
partition 0, leader 2, replicas: 2,1,3, isrs: 2,1,3
partition 1, leader 3, replicas: 3,1,2, isrs: 3,1,2Each group is coordinated by one node, picked by a hash over the live nodes, and the others answer
NOT_COORDINATOR, which every client follows: two consumers of one group started against node 1 and
node 3 join one group and split the partitions between them. Partition leaders are spread over the
nodes so that clients spread too, but a leader here is only an advertisement. Any node serves any
partition, and every voter is a replica. Liveness comes from raft’s own heartbeats, so a node that
stops answering leaves the broker list and the ISR after QUEEN_KAFKA_CLUSTER_TTL_MS (10 s).
Do not put the nodes behind one load-balanced address. A client reconnects to the address it is
given for each broker, so if every node advertises one address, a client sent to its group’s
coordinator lands on another node, is told NOT_COORDINATOR, and loops. A load balancer is fine in
bootstrap.servers, which a client uses once.
Cluster mode refuses transactions: initTransactions() fails at once with
TransactionalIdAuthorizationException. A transaction is staged on the node its records arrive at,
and in a cluster a producer sends its records to each partition’s leader and its commit to the
transaction coordinator, which are different nodes. Transactional applications need a facade
outside cluster mode.
Performance
We measured the facade on 2026-10-01 on the benchmark rig: three 16-vCPU nodes and three load
machines running franz-go with acks=all, idempotence and lz4, on the builds that became
2.0.0-beta.5. At 1,000,000 msg/s into one topic of 200 partitions, records reached their consumers
in 35 ms at the median and 91 ms at p99. Kafka 4.3.1 on the same machines is faster at that shape,
with a p99 of 18 ms. At 100,000 partitions the order reverses: the facade held a p99 of 251 ms at
the same rate, while Kafka fell behind, at 10.9 s. The full matrix, the CPU it cost and what is still
open are on Kafka clients in the benchmarks.
Limits
- Fetch stops at version 6, before fetch sessions (KIP-227), so a consumer names all of its partitions in every fetch. At hundreds of partitions that costs nothing; at 100,000 it is most of the remaining gap to Queen’s own clients, which had a median two to three times lower than the facade’s 134 ms in those runs.
- Records are stored decoded, one message each, and fetch answers carry them in uncompressed batches. A producer’s compression saves bandwidth into the broker only, and librdkafka-based producers do not compress zstd at all here (librdkafka waits for Fetch version 10).
- No KIP-848 groups and no static membership: keep
group.protocol=classic(the default of every Kafka 4.x client) and leavegroup.instance.idunset, or the consumer fails when it joins. - No DeleteRecords, so
kafka-delete-records.shfails and Kafka Streams cannot purge its repartition topics. - On a topic with
retention.msset, a partition that retention has emptied and that nobody has written to for 30 days (PARTITION_CLEANUP_DAYS) is removed, and the next record written to it starts again at offset 0. - Past the cluster’s ceiling (about 1.5M msg/s on the benchmark rig), Kafka producers can keep being admitted while Kafka consumers fall behind. Queen’s own clients kept consuming at the same offered rate, and we have not found the cause yet.
- A queue created by a Queen client shows up in Kafka metadata within about 3 seconds, which is how often the facade re-reads the catalog. A Kafka consumer that subscribed inside that window saw no partitions and read nothing for about 30 seconds in our test, so create the topic first or subscribe after the first push has landed.
- The facade’s own tests run in CI; the client suites under
protocols/queen-kafka/compat/are run by hand. The Jepsen campaign behind guarantees used the HTTP API and did not cover the facade.
Next
- Kafka reference: every setting, the APIs and versions, the topic configs, where the facade differs from Apache Kafka, and what each client library needs.
- Kafka Streams, Connect and Flink: what the frameworks ask of a broker and what runs today.