---
title: "Kafka clients"
description: "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."
---

> Queen MQ documentation, for AI agents
> Complete self-contained summary of Queen MQ: https://queenmq.com/llms-brief.txt
> Fetch that first when the question is about the product rather than about this page.
> Index of all pages: https://queenmq.com/llms.txt

# Kafka clients

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

```bash
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](https://github.com/edenhill/kcat) and read them back:

```bash
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'
```

```text
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:

```bash
curl -s 'localhost:6632/api/v1/pop/queue/orders?batch=10' | jq -c '.partition, .messages[].data'
```

```text
"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:

```bash
kcat -b localhost:9092 -G billing -o beginning -q -f 'billing [%p] offset %o key=%k\n' orders
```

```text
billing [400] offset 0 key=customer-456
billing [307] offset 0 key=customer-123
billing [307] offset 1 key=customer-123
```

```bash
curl -s localhost:6632/api/v1/consumer-groups \
  | jq -c '.[] | select(.name == "billing") | {name, queueName, partitionCursors, totalLag}'
```

```text
{"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`:

```bash
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}'
```

```text
{"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:

```bash
kcat -C -b localhost:9092 -t orders -p 307 -o 2 -c 1 -q -f 'offset %o key=%k value=%s\n'
```

```text
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.)

**Figure.** 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.

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.

- Kafka client: franz-go, librdkafka, Java
- Queen client: six SDKs, or curl
- Kafka facade: port 9092, in the broker
- HTTP edge: port 6632
- one pipeline: admission, raft log, reads
- queue orders: a partition named 307
- Kafka client → Kafka facade: produce, fetch
- Queen client → HTTP edge: push, pop
- Kafka facade → one pipeline: in memory
- HTTP edge → one pipeline
- one pipeline → queue orders: one entry per push

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](#idempotent-producers). 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
`transactionId`s, repeats are stored twice until you give the queue a window
(`dedupWindowSeconds`, see [queue options](/concepts/partitions/#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](/reference/kafka/)). 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](/guides/kafka-frameworks/)
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](/operate/cluster/)), this override sits next to
`compose.yaml` in the same folder:

```yaml title="compose.kafka.yaml"
# 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:

```bash
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
```

```text
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](/benchmarks/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](/concepts/guarantees/) used the HTTP API and did
  not cover the facade.

## Next

- [Kafka reference](/reference/kafka/): 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](/guides/kafka-frameworks/): what the frameworks ask of a broker
  and what runs today.

Source: https://queenmq.com/guides/kafka/index.mdx
