sondahub

sondahub / Kafka test broker

Kafka test broker

A Kafka cluster with something to say: 16 topics carrying the hub’s live activity — orders, stock, IoT telemetry, exchange rates, posts, tickets, flights, sign-ins — with 5 minutes of history in them, consumer groups that rebalance, SASL, and an order topic that answers what you produce. The Kafka protocol itself, through a bridge that lends the broker a port on your computer.

Free, nothing to install but the bridge. Each run of the bridge is a cluster of its own — run several side by side.

Connect

Address
wss://api.sondahub.com/kafka
From a Kafka client
through the bridge: node bridge.mjs kafka [port]
bootstrap.servers
localhost:9192, or the port the bridge printed
Cluster
one broker, node 1; replication factor 1
Security
PLAINTEXT, or SASL_PLAINTEXT with PLAIN, SCRAM-SHA-256, SCRAM-SHA-512
Text messages
JSON, for any WebSocket client

The bridge listens on port 9192 — not Kafka’s usual 9092, so a broker you already run keeps it — or the next free port, and carries each connection a client opens to the hub over one WebSocket. The hub answers as a Kafka broker, advertising the address and port the client dialed, and the bridge prints what it sees: clients, groups and their rebalances, what is produced and read.

Run the bridge
curl -O https://sondahub.com/industrial/bridge.mjs
node bridge.mjs kafka

Kafka on localhost:9192 → wss://api.sondahub.com/kafka
bootstrap.servers=localhost:9192 — every run of the bridge is a cluster of its own
the hub is Kafka cluster TBw08p3TqOS9wY3pmh0INg, broker 1, on port 9192 — 16 topics with the last five minutes in them, store.orders.dlq for what is not an order, and any topic you create
client demo-producer connected from 127.0.0.1:53390
store.orders ← order my-order-1 for 42.5: paid in 2 s, shipped in 5, delivered in 9, each a record on store.orders keyed my-order-1
group demo: rebalancing — demo-consumer joined
group demo: generation 1, 1 member, assignor range, leader demo-consumer
fleet.telemetry → demo-consumer reads from 0@0, 1@0, 2@0, 3@0, 4@0, 5@0
Two clusters side by side
node bridge.mjs kafka          # in one terminal
Kafka on localhost:9192 → wss://api.sondahub.com/kafka

node bridge.mjs kafka          # in another
Kafka on localhost:9193 (9192 was taken) → wss://api.sondahub.com/kafka

In a Kafka client

Start the bridge, then point the client at localhost:9192. Nothing else is needed: no credentials unless you want SASL, no topic to create — 16 are already talking, with 5 minutes of history in them for a consumer that starts from the earliest offset.

kcat
kcat -b localhost:9192 -L                                     # the broker and the topics
kcat -b localhost:9192 -C -t fleet.telemetry -o -3 -e         # the last three records of each partition
echo '{"id":"my-order-1","customer_id":7,"total":42.5}' | kcat -b localhost:9192 -P -t store.orders -k my-order-1
kcat -b localhost:9192 -C -t store.orders -o end -f '%k %s\n'   # and watch it get paid, shipped, delivered
Python (confluent-kafka)
from confluent_kafka import Consumer

c = Consumer({'bootstrap.servers': 'localhost:9192', 'group.id': 'demo', 'auto.offset.reset': 'earliest'})
c.subscribe(['fleet.telemetry'])
while True:
    m = c.poll(1.0)
    if m is not None and m.error() is None:
        print(m.key(), m.value())
Go (franz-go)
cl, err := kgo.NewClient(
    kgo.SeedBrokers("localhost:9192"),
    kgo.ConsumerGroup("demo"),
    kgo.ConsumeTopics("bank.fx"),
)
if err != nil {
    panic(err)
}
for {
    cl.PollFetches(context.Background()).EachRecord(func(r *kgo.Record) {
        fmt.Println(string(r.Key), string(r.Value))
    })
}
With SASL (Java client properties)
bootstrap.servers=localhost:9192
security.protocol=SASL_PLAINTEXT
sasl.mechanism=SCRAM-SHA-256
sasl.jaas.config=org.apache.kafka.common.security.scram.ScramLoginModule required username="sonda" password="sondahub";

In LockFlare Sonda: point a Kafka item at localhost:9192 with the bridge running.

Topics

The same activity the hub’s live streams carry, as keyed JSON records: each value a JSON object, each record with the headers content-type: application/json and source: sondahub, each key on the partition the Java client’s default partitioner (murmur2) would give it. Several records a second on bank.fx, more than one on fleet.telemetry, and one every few seconds on the rest.

TopicWhat it carries
store.orders
3 partitions · key: order id
Orders moving through their life — paid, shipped with a tracking number, delivered — one record per change. Produce your own order here and watch it move.
store.inventory
3 partitions · key: SKU
Stock moving in the warehouses: picks and restocks.
fleet.telemetry
6 partitions · key: device serial
Readings from the IoT fleet, about one a second, with metrics that walk instead of jumping.
fleet.devices
3 partitions · key: device serial
Devices going degraded or offline and coming back.
fleet.alerts
3 partitions · key: device serial
Alerts opened and cleared: thresholds, low battery, weak signal, offline.
bank.fx
3 partitions · key: currency pair
Exchange rates with bid and ask, three a second.
bank.transactions
3 partitions · key: account number
Card transactions on checking accounts, posted or pending.
social.posts
3 partitions · key: username
New posts, with their hashtags.
social.likes
3 partitions · key: post id
Likes landing on posts, with the running count.
social.comments
3 partitions · key: post id
Comments on posts.
helpdesk.tickets
3 partitions · key: ticket number
Tickets changing status and being assigned.
helpdesk.messages
3 partitions · key: ticket number
Messages on tickets, from agents and customers.
flights.board
3 partitions · key: flight number
The departures board: boarding, departed, in the air, landed — and delayed or cancelled.
flights.bookings
3 partitions · key: flight number
Bookings and check-ins, with the seat and the seats left.
identity.logins
3 partitions · key: user name
Sign-ins: the method, and why the ones that failed did.
identity.directory
3 partitions · key: user name
People joining and leaving groups, deactivated, reactivated, retitled.
store.orders.dlq
1 partition · key: as produced
Records produced to store.orders that are not an order the hub can take, as they came, with the reason in an error header and where they were in original.topic, original.partition and original.offset.

Any other topic is yours: create it, or produce to it and let auto-creation make it with one partition. Up to 200 topics in a cluster.

Orders that move

store.orders listens. Produce an order to it — a JSON object with customer_id and a total above 0, and an id if you want to pick it — and the hub takes it through its life: paid two seconds later, shipped with a tracking number at five, delivered at nine, each a record on store.orders keyed by the order id. An order over 10,000 is cancelled instead: the payment is declined.

Anything else produced there — not JSON, not an object, no customer, no total — goes to store.orders.dlq as it came, with an error header saying why and original.topic, original.partition and original.offset saying where it was: a dead-letter topic to point a consumer at.

An order, then its life
key my-order-1  {"id":"my-order-1","customer_id":7,"total":42.5}
key my-order-1  {"id":"my-order-1","number":"SO-my-order-1","customer_id":7,"total":42.5,"status":"paid","previous":"pending"}
key my-order-1  {"id":"my-order-1","number":"SO-my-order-1","customer_id":7,"total":42.5,"status":"shipped","previous":"paid","tracking_number":"1Z…"}
key my-order-1  {"id":"my-order-1","number":"SO-my-order-1","customer_id":7,"total":42.5,"status":"delivered","previous":"shipped"}

Consumer groups

The classic group protocol, as a broker runs it: the first member waits 500 ms for others before the group’s first generation; a new member makes everyone rejoin, and heartbeats say so with REBALANCE_IN_PROGRESS; the leader computes the assignment with whatever assignor the members share — range, round robin, sticky, cooperative-sticky, your own — and the broker hands it out; a member that stops heartbeating is dropped when its session times out (6,000 ms to 1,800,000 ms). Committed offsets last as long as the cluster. Start two consumers in one group and watch the partitions split between them in the bridge’s log.

SASL

Optional, on the same port: a client configured for SASL_PLAINTEXT signs in with PLAIN, SCRAM-SHA-256, SCRAM-SHA-512; one configured for PLAINTEXT just connects.

UserMay
sonda
password sondahub
Everything: produce, consume, create and delete topics, groups.
viewer
password viewer
Read and describe: consume, list, describe. Producing, creating and deleting are refused with an authorization error.
No SASL
PLAINTEXT
Everything, as sonda. Authenticating is up to the client: one port takes both kinds of configuration.
Anyone else
or a wrong password
Refused with SASL_AUTHENTICATION_FAILED, and the connection closes.

What the broker answers

The APIs and versions it lists in ApiVersions. Versions from the first flexible one on use the compact encodings and tagged fields. Produce is listed from v0, as Apache Kafka 4 lists it, but answered from v3, the first with record batches.

APIHere
Produce
key 0 · v3–v11 · flexible from v9
Record batches (message format v2) with their CRC-32C checked; acks 0, 1 and all; idempotent producers, with sequence numbers checked and a batch sent twice answered with its first offset. Batches compressed with gzip, snappy, lz4 or zstd are kept as they came.
Fetch
key 1 · v4–v12 · flexible from v12
Waits up to maxWaitMs for minBytes and answers the moment data lands; maxBytes and partitionMaxBytes as KIP-74 has them. Every fetch is a full one (no fetch sessions).
ListOffsets
key 2 · v1–v7 · flexible from v6
Earliest, latest, the largest timestamp, and the first offset at or after a time.
Metadata
key 3 · v0–v12 · flexible from v9
One broker, node 1, at the address and port the client dialed. A topic asked for that does not exist is created (auto.create.topics.enable), unless the client says not to.
OffsetCommit
key 8 · v2–v9 · flexible from v8
Committed offsets per group, checked against the generation and member; simple commits from outside a group too.
OffsetFetch
key 9 · v1–v9 · flexible from v6
One group, or several at once (v8 on).
FindCoordinator
key 10 · v0–v5 · flexible from v3
Broker 1 coordinates every group. There is no transaction coordinator.
JoinGroup
key 11 · v0–v9 · flexible from v6
The classic group protocol: MEMBER_ID_REQUIRED for a new member, 500 ms for more members before a new group’s first generation, static membership with group.instance.id.
Heartbeat
key 12 · v0–v4 · flexible from v4
REBALANCE_IN_PROGRESS when it is time to rejoin; a member that stops is dropped when its session times out.
LeaveGroup
key 13 · v0–v5 · flexible from v4
One member, or several by member or instance id (v3 on).
SyncGroup
key 14 · v0–v5 · flexible from v4
The leader’s assignments, handed to every member; followers wait for them as they would on a broker.
DescribeGroups
key 15 · v0–v5 · flexible from v5
State, assignor, members with their metadata and assignments.
ListGroups
key 16 · v0–v5 · flexible from v3
With state and type filters.
SaslHandshake
key 17 · v1–v1
PLAIN, SCRAM-SHA-256, SCRAM-SHA-512.
ApiVersions
key 18 · v0–v4 · flexible from v3
Also the KIP-511 answer, in v0, to a version it does not know.
CreateTopics
key 19 · v2–v7 · flexible from v5
Up to 64 partitions, replication factor 1, the topic configurations below; validateOnly too.
DeleteTopics
key 20 · v1–v5 · flexible from v4
Your topics; the hub’s stay.
InitProducerId
key 22 · v0–v5 · flexible from v2
Producer ids and epochs for idempotent producers. A transactional id is refused: there are no transactions.
OffsetForLeaderEpoch
key 23 · v2–v4 · flexible from v4
Every partition has had one leader, at epoch 0.
DescribeConfigs
key 32 · v1–v4 · flexible from v4
Topics and the broker.
SaslAuthenticate
key 36 · v0–v2 · flexible from v2
The users below; a failed sign-in closes the connection, as a broker does.
CreatePartitions
key 37 · v0–v3 · flexible from v2
More partitions for your topics.
DeleteGroups
key 42 · v0–v2 · flexible from v2
Groups with no members left.
DescribeCluster
key 60 · v0–v1 · flexible from v0
The cluster id, the controller and the one broker.

Configurations and limits

Topic configurationNotes
cleanup.policy
default delete
delete or compact are taken; the log is not compacted here, only kept by retention.
compression.type
default producer
Batches are kept as the producer compressed them.
retention.ms
default 3600000
Batches older than this go. The cluster’s own cap can drop them sooner.
retention.bytes
default -1
A partition bigger than this drops its oldest batches; -1 for no limit.
max.message.bytes
default 1048588
The largest record batch a producer may send to the topic.
message.timestamp.type
default CreateTime
CreateTime keeps the producer’s timestamps; LogAppendTime stamps each batch with the broker’s clock when it lands.
min.insync.replicas
default 1
Accepted and shown back; it changes nothing here.
segment.bytes
default 1073741824
Accepted and shown back; it changes nothing here.
segment.ms
default 604800000
Accepted and shown back; it changes nothing here.
delete.retention.ms
default 86400000
Accepted and shown back; it changes nothing here.
min.compaction.lag.ms
default 0
Accepted and shown back; it changes nothing here.
unclean.leader.election.enable
default false
Accepted and shown back; it changes nothing here.

A cluster keeps up to 24 MB of records; past that the oldest batches go, and the log start offset moves, so a consumer that asks for an offset that is gone gets OFFSET_OUT_OF_RANGE and its reset policy decides. Requests up to 8 MB, record batches up to 1,048,588 bytes (max.message.bytes), 64 partitions a topic. Batches compressed with gzip, snappy and lz4 are opened (so an order compressed any of those ways is taken); zstd batches are kept and served as they came, but not opened.

JSON, from any WebSocket client

A text message to wss://api.sondahub.com/kafka is a request: list the topics, subscribe to records as they land (with what the topics already hold first, if you ask), or produce one. Values and keys come back as JSON when they are, as text when they are not, and in base64 when they are neither.

Requests
{"topics":true}
{"subscribe":"fleet.telemetry"}
{"subscribe":["store.orders","store.orders.dlq"],"from":"earliest"}
{"produce":"store.orders","key":"my-order-2","value":{"customer_id":7,"total":19.9}}
{"unsubscribe":"*"}
A record
{"topic":"bank.fx","partition":1,"offset":193,"timestamp":"2026-10-07T23:14:45.261Z","key":"USD/AUD",
 "value":{"pair":"USD/AUD","base":"USD","quote":"AUD","rate":1.5255,"bid":1.5243,"ask":1.5267,"change_pct":0.36},
 "headers":{"content-type":"application/json","source":"sondahub"}}

The bridge, on the wire

For anyone writing a bridge of their own: it opens wss://api.sondahub.com/kafka with the subprotocol kafka and says hello in text — {"bridge":"kafka","version":1,"port":9192}, the port it listens on, which the broker advertises. Then every client connection rides the WebSocket under an id the bridge gives it: a byte for what happened (3 bytes of the stream, 4 a connection opened, 5 one closed), the id (four bytes, big-endian), two zero bytes, then for a 3 the bytes as they came off the socket, and for a 4 {"peer":"ip:port","local":"ip:port"} in UTF-8 — who dialed and the address they dialed. The hub sends 3 and 5 back, and {"log":…} lines in text.

Questions

Why does Kafka need the bridge?

Kafka clients open TCP connections, and only HTTP reaches sondahub — no raw TCP comes in. So the broker runs on the hub and the bridge lends it a port on your computer, carrying every connection a client opens — the bootstrap one, the broker one, the group coordinator’s — over one WebSocket. To the client it is a broker on localhost.

Why port 9192 and not 9092?

So a real broker on your computer keeps 9092, and so you can run several: the bridge takes 9192, or the next free port up to 9201, and prints the one it took. Run it twice and you have two clusters side by side, each with its own topics, groups and offsets. Give it a port — node bridge.mjs kafka 9092 — and it takes exactly that one or stops.

Do my records reach other people?

No. Each run of the bridge is a cluster of its own: the clients you point at it share it — produce from one, consume from another — but nobody else sees it, and it ends when the bridge stops. When the bridge reconnects, the hub is a new cluster, and the clients start over against it.

Which clients work?

Any that speaks the Kafka protocol to a broker of Apache Kafka 2.1 or newer: the broker answers the APIs and versions in the table below, and a client picks the highest version both sides know from ApiVersions. The message formats come from the message definitions in Apache Kafka’s own source. The new consumer group protocol (KIP-848) is not offered, so a consumer uses the classic one — the default.

Are there transactions, TLS or a Schema Registry?

No. A transactional producer is refused when it asks for a producer id; the broker speaks plaintext, with SASL or without, and no TLS; there is no Schema Registry. Records are bytes: the hub’s own are JSON.