---
title: "Kafka bridge"
description: "A Kafka producer writes orders, a Queen worker bills them with one transaction per order, and a Kafka consumer reads the invoices, all on the same queues of one node, with no connector."
---

> 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 bridge

You do not have to move a Kafka shop to Queen all at once. In this program an order service keeps
its Kafka producer and changes only `bootstrap.servers`, an analytics service keeps its Kafka
consumer, and the one piece in between, billing, is written against Queen so it can commit each
order's ack, its invoice and the customer's running total as one transaction. Nothing copies data
between the two worlds: a Kafka topic is a Queen queue, and Kafka partition 3 is the Queen
partition named `"3"`, so both kinds of client read and write the same messages.

## Run it

The node needs its Kafka listener on:

```bash
docker run --platform linux/amd64 -d --name queen -p 6632:6632 -p 9092:9092 \
  -e QUEEN_KAFKA_EMBEDDED=true \
  -e QUEEN_KAFKA_ADVERTISED_ADDR=localhost:9092 \
  ghcr.io/queen-mq/queen:latest
npm install queen-mq @confluentinc/kafka-javascript
node bridge.mjs     # the program below, saved as bridge.mjs
```

What it printed against a 2.0.0-beta.6 node, with `@confluentinc/kafka-javascript` 1.10.1:

```text
producing orders with a Kafka client
  12 orders for 3 customers

billing with a Queen worker
  stored payload: {"h":[{"k":"source","v":"c2hvcA=="}],"k":"YWxpY2U=","t":1790948491582,"v":"eyJjdXN0b21lciI6ImFsaWNlIiwibiI6MSwiY2VudHMiOjEwMDB9"}
  ok: alice's orders arrived in order, all from Kafka partition 1
  ok: bob's orders arrived in order, all from Kafka partition 2
  ok: carol's orders arrived in order, all from Kafka partition 2
  ok: billing saw all 12 orders with their Kafka key and header

reading invoices with a Kafka consumer
  first: partition 2, key null, value {"cents":1000,"customer":"bob","invoice":"INV-bob-1","n":1}
  ok: analytics read all 12 invoices (got 12)
  ok: each invoice is in the partition number its order came from

  totals: alice=10000, bob=10000, carol=10000
  ok: each customer's total is 10000, counted once per order

PASS: 7 checks
```

The `stored payload` line is a Kafka record as Queen keeps it: key, value and header values in
base64, because Kafka bytes are not always text, and the producer's timestamp in `t`. The last
lines go the other way: a payload a Queen client pushed reaches a Kafka consumer as its JSON
text, with no key.

## The program

```js title="examples/cross-protocol/bridge.mjs"
//
// Kafka producers in, a Queen worker in the middle, a Kafka consumer out, all
// on one broker and the same queues.
//
// An order service that already produces to Kafka keeps its Kafka client and
// changes only bootstrap.servers. Billing is written against Queen: for each
// order it commits the ack, an invoice and the customer's running total as one
// transaction. Analytics keeps its Kafka consumer and reads the invoices. There
// is no connector and no second copy of anything: a Kafka topic is a Queen
// queue, and Kafka partition n is the Queen partition named "n".
//
//   order service (Kafka producer) -> orders            (Kafka topic = Queen queue)
//   billing (Queen worker)         -> ack + invoices + KV, one transaction
//   analytics (Kafka consumer)     <- invoices
//
// Needs a node with the Kafka listener on:
//   docker run -d -p 6632:6632 -p 9092:9092 -e QUEEN_KAFKA_EMBEDDED=true \
//     -e QUEEN_KAFKA_ADVERTISED_ADDR=localhost:9092 ghcr.io/queen-mq/queen:latest
//
// Run it:
//   npm install && node bridge.mjs

import Confluent from '@confluentinc/kafka-javascript'
import { Queen } from 'queen-mq'

const QUEEN_URL = process.env.QUEEN_URL || 'http://localhost:6632'
const KAFKA = process.env.KAFKA || 'localhost:9092'
const RUN = Date.now().toString(36)
const ORDERS = `bridge-orders-${RUN}`
const INVOICES = `bridge-invoices-${RUN}`
const TOTALS = `bridge-totals-${RUN}` // KV namespace: revenue per customer

const CUSTOMERS = ['alice', 'bob', 'carol']
const ORDERS_PER_CUSTOMER = 4
const TOTAL_ORDERS = CUSTOMERS.length * ORDERS_PER_CUSTOMER

let checks = 0
const assert = (condition, description) => {
  if (!condition) throw new Error(description)
  checks++
  console.log(`  ok: ${description}`)
}

const { Kafka, logLevel } = Confluent.KafkaJS
const kafka = new Kafka({ 'bootstrap.servers': KAFKA, kafkaJS: { logLevel: logLevel.WARN } })
const queen = new Queen({ url: QUEEN_URL, handleSignals: false })

// The facade stores a Kafka record as JSON with base64 bytes, because a Kafka
// key or value is arbitrary bytes and a Queen payload is JSON:
//   { k: key, v: value, h: [{ k: name, v: value }], t: timestamp }
// h is left out when the record has no headers, t when it has no timestamp.
const b64 = (s) => (s == null ? null : Buffer.from(s, 'base64').toString())
const decode = (data) => ({
  key: b64(data.k),
  value: b64(data.v),
  headers: Object.fromEntries((data.h ?? []).map(h => [h.k, b64(h.v)])),
})

let analytics = null

try {
  console.log(`queen ${QUEEN_URL}, kafka ${KAFKA}`)

  // Both topics are created up front through the Kafka admin API, as a Kafka
  // deployment would. Four partitions each: the producer hashes customers onto
  // them, so each customer's orders stay in one partition, in order.
  const admin = kafka.admin()
  await admin.connect()
  await admin.createTopics({
    topics: [ORDERS, INVOICES].map(topic => ({ topic, numPartitions: 4 })),
  })
  await admin.disconnect()

  // -------------------------------------------------- the order service, Kafka
  console.log('\nproducing orders with a Kafka client')
  const producer = kafka.producer({ kafkaJS: { acks: -1 } })
  await producer.connect()
  for (let n = 1; n <= ORDERS_PER_CUSTOMER; n++) {
    for (const customer of CUSTOMERS) {
      await producer.send({
        topic: ORDERS,
        messages: [{
          key: customer,
          value: JSON.stringify({ customer, n, cents: 1000 * n }),
          headers: { source: 'shop' },
        }],
      })
    }
  }
  await producer.disconnect()
  console.log(`  ${TOTAL_ORDERS} orders for ${CUSTOMERS.length} customers`)

  // ---------------------------------------------------------- billing, Queen
  //
  // A Queen consumer group on the topic. For each order, one transaction acks
  // it, pushes the invoice into the same partition number of the invoices
  // topic, and adds the order to the customer's total in KV. If the worker dies
  // halfway, none of the three happened and the order comes back.
  console.log('\nbilling with a Queen worker')
  const billed = []
  let shownRaw = false
  await queen
    .queue(ORDERS)
    .group('billing')
    .subscriptionMode('all')
    .autoAck(false) // the ack rides the transaction
    .each()
    .limit(TOTAL_ORDERS)
    .timeoutMillis(1000)
    .idleMillis(10000)
    .consume(async (msg) => {
      if (!shownRaw) {
        console.log(`  stored payload: ${JSON.stringify(msg.data)}`)
        shownRaw = true
      }
      const record = decode(msg.data)
      const order = JSON.parse(record.value)

      await queen
        .transaction()
        .ack(msg, 'completed', { consumerGroup: 'billing' })
        .queue(INVOICES).partition(msg.partition)
        .push({ data: { invoice: `INV-${order.customer}-${order.n}`, ...order } })
        .kv.incr(TOTALS, order.customer, order.cents, { ttl: '1h' })
        .commit()

      billed.push({ partition: msg.partition, key: record.key, source: record.headers.source, ...order })
    })

  for (const customer of CUSTOMERS) {
    const mine = billed.filter(b => b.customer === customer)
    assert(
      mine.map(b => b.n).join(',') === '1,2,3,4' && new Set(mine.map(b => b.partition)).size === 1,
      `${customer}'s orders arrived in order, all from Kafka partition ${mine[0]?.partition}`
    )
  }
  assert(
    billed.length === TOTAL_ORDERS && billed.every(b => b.key === b.customer && b.source === 'shop'),
    `billing saw all ${TOTAL_ORDERS} orders with their Kafka key and header`
  )

  // ------------------------------------------------------- analytics, Kafka
  //
  // A Kafka consumer group on the invoices that Queen wrote. A native payload
  // reaches a Kafka client as its JSON text, with no key.
  console.log('\nreading invoices with a Kafka consumer')
  analytics = kafka.consumer({ kafkaJS: { groupId: `analytics-${RUN}`, fromBeginning: true } })
  await analytics.connect()
  await analytics.subscribe({ topics: [INVOICES] })
  const invoices = []
  await new Promise((resolve, reject) => {
    const deadline = setTimeout(() => resolve(), 30000)
    analytics.run({
      eachMessage: async ({ partition, message }) => {
        invoices.push({ partition, key: message.key, invoice: JSON.parse(message.value.toString()) })
        if (invoices.length === TOTAL_ORDERS) {
          clearTimeout(deadline)
          resolve()
        }
      },
    }).catch(reject)
  })
  console.log(`  first: partition ${invoices[0]?.partition}, key ${invoices[0]?.key}, value ${JSON.stringify(invoices[0]?.invoice)}`)
  assert(invoices.length === TOTAL_ORDERS, `analytics read all ${TOTAL_ORDERS} invoices (got ${invoices.length})`)
  assert(
    invoices.every(i => billed.some(b => b.n === i.invoice.n && b.customer === i.invoice.customer && Number(b.partition) === i.partition)),
    'each invoice is in the partition number its order came from'
  )

  // ------------------------------------------------------------ the totals
  const totals = await queen.kv.getMany(TOTALS, CUSTOMERS)
  const expected = 1000 * (ORDERS_PER_CUSTOMER * (ORDERS_PER_CUSTOMER + 1)) / 2
  console.log(`\n  totals: ${totals.rows.map(r => `${r.key}=${r.value}`).join(', ')}`)
  assert(
    totals.rows.length === CUSTOMERS.length && totals.rows.every(r => Number(r.value) === expected),
    `each customer's total is ${expected}, counted once per order`
  )

  console.log(`\nPASS: ${checks} checks`)
} catch (err) {
  console.error(`\nFAIL: ${err.message}`)
  process.exitCode = 1
} finally {
  try {
    if (analytics) await analytics.disconnect()
    for (const customer of CUSTOMERS) await queen.kv.delete(TOTALS, customer)
    for (const q of [ORDERS, INVOICES]) await queen.queue(q).delete()
  } catch (err) {
    console.error(`  (cleanup incomplete: ${err.message})`)
  }
  await queen.close()
}
```

## How it works

The program creates both topics through the Kafka admin API, four partitions each, as a Kafka
deployment would. The producer sends each order with the customer as its key, and the Kafka
client hashes the key onto one of the four partitions, so each customer's orders sit in one
partition, in order, and two customers can share one (bob and carol did).

**Figure.** The Kafka bridge on one Queen node. An order service with a Kafka producer writes orders, keyed by customer, to the topic orders, which is the Queen queue orders with four partitions. Billing, a Queen consumer group, pops each order and commits one transaction: the ack of the order, the invoice pushed into the partition of the invoices topic with the same number, and a KV incr of the customer's running total. Analytics, a Kafka consumer group, fetches the invoices from the beginning.

There is one copy of each record, so there is no connector to run and no lag between copies to watch.

- order service: Kafka producer
- orders: topic = queue, 4 partitions
- billing: Queen consumer
- one transaction: ack the order push the invoice incr the total
- invoices: same partition number
- analytics: Kafka consumer
- order service → orders: produce
- orders → billing: pop
- billing → one transaction: commit
- one transaction → invoices: push
- invoices → analytics: fetch

Billing is an ordinary Queen consumer group on the `orders` queue. It decodes the record from the
payload and, for each order, commits one transaction: the ack of the order, the invoice pushed into
the partition of the `invoices` topic with the same number as the order's, and a `kv.incr` of the
customer's total. If the worker dies halfway, none of the three happened and the order comes back,
which is why the totals in the run are exact. On Kafka the same guarantee takes a transactional
producer, offsets sent to the transaction and consumers reading committed data only, and the
running total would still live somewhere else.

Analytics is an ordinary Kafka consumer group on `invoices`, reading from the beginning.
[Kafka clients](/guides/kafka/) covers how topics, partitions, groups, offsets and transactions
map onto Queen, and what the listener does and does not support.

## Why it is short

Putting a non-Kafka worker into a Kafka pipeline usually takes a connector in each direction,
which means every record copied out and back again, and a lag between the copies to watch. Here
there is one copy of each record, so the program is the three services and nothing in between.

## Limits

- Kafka decides the partitions. A Kafka producer hashes keys onto a fixed number of numbered
  partitions, so the Queen worker inherits that layout, shared partitions included. Native
  producers can give every customer a partition of its own, but a Kafka client cannot address a
  partition whose name is not a number.
- Create topics before Kafka consumers subscribe. The listener caches the list of queues for
  three seconds, so a queue that a native push has just created is unknown to Kafka metadata until
  the cache refreshes; we measured 3.0 s. A Kafka consumer that subscribed in that window was told
  the topic does not exist and read nothing in the 30 seconds we waited.
- A topic created through Kafka has no dedup window. Kafka records carry no Queen
  `transactionId`, so such a topic gets `dedupWindowSeconds: 0`, and a native push into it with a
  repeated `transactionId` is stored twice; we checked. Set `dedupWindowSeconds` on the queue if
  native producers count on it.
- The program needs a node with the listener on, which is why it lives in
  `examples/cross-protocol/` and not with the other programs that `examples/apps/run.sh` runs.

## Next

[Kafka clients](/guides/kafka/) for the whole mapping and the clients that were tested, and the
[card charger](/examples/exactly-once/) for the transaction that billing uses.

Source: https://queenmq.com/examples/kafka-bridge/index.mdx
