Skip to content

Kafka clients

Connect an unchanged Kafka producer, consumer or admin tool to Queen: how topics, partitions, consumer groups, transactions and compaction map onto queues, what it measured, and where it stops.

Updated View as Markdown

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:latest

QUEEN_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' orders
billing [400] offset 0 key=customer-456
billing [307] offset 0 key=customer-123
billing [307] offset 1 key=customer-123
curl -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.)

Kafka clients connect to the Kafka facade on port 9092, which runs inside the Queen broker process. Queen's own clients connect to the HTTP edge on port 6632. A produce becomes a push and a fetch a read, called in memory, and both paths go through the same admission, the same raft log and the same reads. Topic orders is the queue orders, and Kafka partition 307 is the Queen partition named 307.Kafka clientfranz-go, librdkafka, JavaQueen clientsix SDKs, or curlKafka facadeport 9092, in the brokerHTTP edgeport 6632one pipelineadmission, raft log, readsqueue ordersa partition named 307produce, fetchpush, popin memoryone entry per push
One pipeline behind two protocols. A Kafka record is a Queen message from the moment it is written, so replication, retention, the dashboard and Queen's own consumers all apply to it. Source: server/src/kafka_inproc.rs, protocols/queen-kafka/src/handlers/produce.rs

With 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:

compose.kafka.yamlyaml
# 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 -9
Metadata 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,2

Each 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 leave group.instance.id unset, or the consumer fails when it joins.
  • No DeleteRecords, so kafka-delete-records.sh fails and Kafka Streams cannot purge its repartition topics.
  • On a topic with retention.ms set, 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.
Navigation

Type to search…

↑↓ navigate↵ selectEsc close