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.
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
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 -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
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())
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))
})
}
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.
| Topic | What 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.
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.
| User | May |
|---|---|
| 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.
| API | Here |
|---|---|
| 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 configuration | Notes |
|---|---|
| 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.
{"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":"*"}
{"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.