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:
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.mjsWhat it printed against a 2.0.0-beta.6 node, with @confluentinc/kafka-javascript 1.10.1:
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 checksThe 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
//
// 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).
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 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 getsdedupWindowSeconds: 0, and a native push into it with a repeatedtransactionIdis stored twice; we checked. SetdedupWindowSecondson 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 thatexamples/apps/run.shruns.
Next
Kafka clients for the whole mapping and the clients that were tested, and the card charger for the transaction that billing uses.