KubeMQ
ConnectorsKafkaTutorials

Getting Started

Start kubemq-server, then produce and consume your first message over the Kafka wire protocol with kcat and seven client libraries — a full round-trip.

Get a message flowing through the KubeMQ Kafka connector in minutes. You point a stock Kafka client — kcat, or any of seven client libraries — at KubeMQ by repointing bootstrap.servers, produce one record to a topic, and consume it back with a consumer group. There's no KubeMQ SDK and no client-library swap: the connector speaks the real Produce/Fetch/group-coordinator wire protocol, so any unmodified Kafka client just works. By the end of this page you'll have run a complete produce-then-consume round-trip against a local KubeMQ server, using whichever client you already have installed.

Prerequisites

  • A running kubemq-server, reachable on port 9092 (plain TCP). The Kafka connector is on by default, so a stock server already serves it — the step below shows how to start one.
  • One of the eight clients in the tabs below: kcat, or a client library for your language. There is no KubeMQ SDK — every example on this page is a stock, unmodified Kafka client.

Start the server

The Kafka connector is on by default — a stock kubemq-server binds port 9092 as soon as it starts, with nothing to enable. Run one locally:

docker run -d \  --pull always \  --platform linux/amd64 \  --name kubemq \  --hostname kubemq \  -p 127.0.0.1:9092:9092 \  -p 127.0.0.1:9093:9093 \  -p 127.0.0.1:50000:50000 \  -p 127.0.0.1:8080:8080 \  -e STORE_ENGINE=next \  -e STORE_NEXT_ACK_POLICY=strict \  -e STORE_STORE_PATH=/kubemq/store \  -e API_BIND_ADDRESS=0.0.0.0 \  -v kubemq-data:/kubemq/store \  europe-docker.pkg.dev/kubemq/images/kubemq-next:latest

Don't want Kafka on this server? Set CONNECTORS_KAFKA_ENABLE=false. For Kubernetes, set spec.kafka.enabled: false in the KubemqCluster CR (Helm: kafka.enabled: false in your values file) — omitting the field means on.

Port 9092 already in use? If a real Kafka broker (or anything else) already holds 9092 on this host, kubemq-server logs an error and keeps running without the Kafka connector — the dashboard still lists it as enabled. Stop the other process, or move KubeMQ's listener with CONNECTORS_KAFKA_PORT and point bootstrap.servers at the new port.

Connecting from another machine? A Kafka client follows the broker address returned after its first connection. Set CONNECTORS_KAFKA_ADVERTISED_HOST to an address that client can reach, or it can connect and then stall. Configure remote or Kubernetes access.

At startup the connector binds two listeners — 9092 for plain TCP and 9093 for TLS (once TLS is configured) — and every produced record lands on KubeMQ's auto-selected next storage engine. You don't need to configure this yourself: on a fresh store, Kafka auto-selects next with no manual store.engine step. See Storage Engines → Zero-config engine selection for the full selection rules, including what happens on an existing store.

You also don't need to pre-create the orders topic used below — the connector creates a topic the first time a producer writes to it, the same as real Kafka's auto.create.topics.enable default. An auto-created topic gets one partition, as Apache Kafka's num.partitions default does; that count is not configurable here, so create the topic explicitly (kafka-topics --create, or an AdminClient) when you want more. A consumer that starts before the topic exists does not create it: a client that asked not to auto-create is told the topic is unknown, and waits for it, exactly as it would against Kafka.

What differs from Kafka

Before moving a production workload, review the behaviors that differ from Kafka: acks=0 on a follower, retention by size, and client-side key partitioning.

Produce a message

Every example below produces one record with the value hello kubemq to the topic orders on bootstrap.servers=localhost:9092 (or, for kcat, -b localhost:9092). No KubeMQ SDK — just a stock Kafka client, repointed.

kcat (the librdkafka CLI) needs no client code at all — pipe the payload straight to the broker:

echo "hello kubemq" | kcat -b localhost:9092 -t orders -P

The Go example uses franz-go, the client this connector's own conformance harness is validated against:

package main

import (
	"context"
	"log"

	"github.com/twmb/franz-go/pkg/kgo"
)

func main() {
	client, err := kgo.NewClient(
		kgo.SeedBrokers("localhost:9092"),
		kgo.DefaultProduceTopic("orders"),
	)
	if err != nil {
		log.Fatalf("client: %v", err)
	}
	defer client.Close()

	record := &kgo.Record{Value: []byte("hello kubemq")}
	if err := client.ProduceSync(context.Background(), record).FirstErr(); err != nil {
		log.Fatalf("produce: %v", err)
	}
	log.Println("produced: hello kubemq")
}
from confluent_kafka import Producer

producer = Producer({"bootstrap.servers": "localhost:9092"})


def delivery_report(err, msg):
    if err is not None:
        raise RuntimeError(f"delivery failed: {err}")
    print(f"produced: {msg.value().decode()} (partition {msg.partition()}, offset {msg.offset()})")


producer.produce("orders", value=b"hello kubemq", callback=delivery_report)
producer.flush(10)
import org.apache.kafka.clients.producer.KafkaProducer;
import org.apache.kafka.clients.producer.ProducerRecord;
import org.apache.kafka.clients.producer.RecordMetadata;
import org.apache.kafka.common.serialization.StringSerializer;

import java.util.Properties;
import java.util.concurrent.ExecutionException;

public final class Produce {
    public static void main(String[] args) throws ExecutionException, InterruptedException {
        Properties props = new Properties();
        props.put("bootstrap.servers", "localhost:9092");
        props.put("key.serializer", StringSerializer.class.getName());
        props.put("value.serializer", StringSerializer.class.getName());

        try (KafkaProducer<String, String> producer = new KafkaProducer<>(props)) {
            RecordMetadata meta = producer.send(new ProducerRecord<>("orders", "hello kubemq")).get();
            System.out.printf("produced: hello kubemq (partition %d, offset %d)%n",
                    meta.partition(), meta.offset());
        }
    }
}
const { Kafka } = require("kafkajs");

const kafka = new Kafka({ brokers: ["localhost:9092"] });
const producer = kafka.producer();

async function main() {
  await producer.connect();
  await producer.send({
    topic: "orders",
    messages: [{ value: "hello kubemq" }],
  });
  console.log("produced: hello kubemq");
  await producer.disconnect();
}

main().catch((err) => {
  console.error("produce failed:", err);
  process.exit(1);
});
using Confluent.Kafka;

var config = new ProducerConfig { BootstrapServers = "localhost:9092" };

using var producer = new ProducerBuilder<Null, string>(config).Build();

var result = await producer.ProduceAsync("orders", new Message<Null, string> { Value = "hello kubemq" });
Console.WriteLine($"produced: hello kubemq (partition {result.Partition}, offset {result.Offset})");
require "rdkafka"

config = Rdkafka::Config.new("bootstrap.servers" => "localhost:9092")
producer = config.producer

handle = producer.produce(topic: "orders", payload: "hello kubemq")
handle.wait(max_wait_timeout_ms: 10_000)
puts "produced: hello kubemq"

producer.close
use rdkafka::config::ClientConfig;
use rdkafka::producer::{FutureProducer, FutureRecord};
use std::time::Duration;

#[tokio::main]
async fn main() {
    let producer: FutureProducer = ClientConfig::new()
        .set("bootstrap.servers", "localhost:9092")
        .create()
        .expect("producer creation failed");

    let record = FutureRecord::<str, _>::to("orders").payload("hello kubemq");
    match producer.send(record, Duration::from_secs(10)).await {
        Ok((partition, offset)) => {
            println!("produced: hello kubemq (partition {partition}, offset {offset})")
        }
        Err((err, _)) => eprintln!("produce failed: {err}"),
    }
}

Whichever client you ran, the connector auto-created the orders topic on that first write and appended your record to it. The next step reads it back.

Consume and verify

Consume the record back with a consumer group named orders-group. Each example joins the group, reads one record, and prints its value alongside the partition and offset the connector assigned — the offset maps one-to-one onto the underlying Events Store Sequence, durable and stable across a restart.

-G puts kcat into consumer-group mode; -c 1 exits after one message:

kcat -b localhost:9092 -G orders-group -c 1 orders
package main

import (
	"context"
	"fmt"
	"log"
	"time"

	"github.com/twmb/franz-go/pkg/kgo"
)

func main() {
	client, err := kgo.NewClient(
		kgo.SeedBrokers("localhost:9092"),
		kgo.ConsumerGroup("orders-group"),
		kgo.ConsumeTopics("orders"),
	)
	if err != nil {
		log.Fatalf("client: %v", err)
	}
	defer client.Close()

	ctx, cancel := context.WithTimeout(context.Background(), 10*time.Second)
	defer cancel()

	fetches := client.PollFetches(ctx)
	if errs := fetches.Errors(); len(errs) > 0 {
		log.Fatalf("fetch: %v", errs)
	}
	fetches.EachRecord(func(record *kgo.Record) {
		fmt.Printf("consumed: %s (partition %d, offset %d)\n", record.Value, record.Partition, record.Offset)
	})
}
from confluent_kafka import Consumer

consumer = Consumer({
    "bootstrap.servers": "localhost:9092",
    "group.id": "orders-group",
    "auto.offset.reset": "earliest",
})
consumer.subscribe(["orders"])

msg = consumer.poll(10.0)
if msg is None:
    raise SystemExit("no message received within timeout")
if msg.error():
    raise RuntimeError(msg.error())

print(f"consumed: {msg.value().decode()} (partition {msg.partition()}, offset {msg.offset()})")
consumer.close()
import org.apache.kafka.clients.consumer.ConsumerRecord;
import org.apache.kafka.clients.consumer.ConsumerRecords;
import org.apache.kafka.clients.consumer.KafkaConsumer;
import org.apache.kafka.common.serialization.StringDeserializer;

import java.time.Duration;
import java.util.List;
import java.util.Properties;

public final class Consume {
    public static void main(String[] args) {
        Properties props = new Properties();
        props.put("bootstrap.servers", "localhost:9092");
        props.put("group.id", "orders-group");
        props.put("key.deserializer", StringDeserializer.class.getName());
        props.put("value.deserializer", StringDeserializer.class.getName());
        props.put("auto.offset.reset", "earliest");

        try (KafkaConsumer<String, String> consumer = new KafkaConsumer<>(props)) {
            consumer.subscribe(List.of("orders"));
            ConsumerRecords<String, String> records = consumer.poll(Duration.ofSeconds(10));
            for (ConsumerRecord<String, String> record : records) {
                System.out.printf("consumed: %s (partition %d, offset %d)%n",
                        record.value(), record.partition(), record.offset());
            }
        }
    }
}
const { Kafka } = require("kafkajs");

const kafka = new Kafka({ brokers: ["localhost:9092"] });
const consumer = kafka.consumer({ groupId: "orders-group" });

async function main() {
  await consumer.connect();
  await consumer.subscribe({ topic: "orders", fromBeginning: true });

  await consumer.run({
    eachMessage: async ({ partition, message }) => {
      console.log(`consumed: ${message.value.toString()} (partition ${partition}, offset ${message.offset})`);
      await consumer.disconnect();
    },
  });
}

main().catch((err) => {
  console.error("consume failed:", err);
  process.exit(1);
});
using Confluent.Kafka;

var config = new ConsumerConfig
{
    BootstrapServers = "localhost:9092",
    GroupId = "orders-group",
    AutoOffsetReset = AutoOffsetReset.Earliest,
};

using var consumer = new ConsumerBuilder<Ignore, string>(config).Build();
consumer.Subscribe("orders");

var result = consumer.Consume(TimeSpan.FromSeconds(10));
if (result is null)
{
    throw new TimeoutException("no message received within timeout");
}
Console.WriteLine($"consumed: {result.Message.Value} (partition {result.Partition}, offset {result.Offset})");

consumer.Close();
require "rdkafka"

config = Rdkafka::Config.new(
  "bootstrap.servers" => "localhost:9092",
  "group.id" => "orders-group",
  "auto.offset.reset" => "earliest"
)
consumer = config.consumer
consumer.subscribe("orders")

message = consumer.poll(10_000)
raise "no message received within timeout" if message.nil?

puts "consumed: #{message.payload} (partition #{message.partition}, offset #{message.offset})"

consumer.close
use rdkafka::config::ClientConfig;
use rdkafka::consumer::{BaseConsumer, Consumer};
use rdkafka::message::Message;
use std::time::Duration;

fn main() {
    let consumer: BaseConsumer = ClientConfig::new()
        .set("bootstrap.servers", "localhost:9092")
        .set("group.id", "orders-group")
        .set("auto.offset.reset", "earliest")
        .create()
        .expect("consumer creation failed");

    consumer.subscribe(&["orders"]).expect("subscribe failed");

    match consumer.poll(Duration::from_secs(10)) {
        Some(Ok(message)) => {
            let payload = message
                .payload()
                .map(|p| String::from_utf8_lossy(p).to_string())
                .unwrap_or_default();
            println!("consumed: {payload} (partition {}, offset {})", message.partition(), message.offset());
        }
        Some(Err(err)) => eprintln!("consume failed: {err}"),
        None => eprintln!("no message received within timeout"),
    }
}

A successful round-trip prints the record you produced, plus the partition and offset the connector assigned it:

produced: hello kubemq
consumed: hello kubemq (partition 0, offset 0)

You just repointed a stock Kafka client at KubeMQ, produced a record, and consumed it back through a real consumer group — the same round-trip you'd run against any Kafka cluster, with no client-library swap and no code change beyond bootstrap.servers. From here, dig into a single feature end-to-end — producing with keys and durability guarantees, consuming with manual offset control, or how topics and partitions map onto KubeMQ's storage layer.

Next steps

Was this page helpful?

On this page