Migrate from Kafka
Try KubeMQ as your Kafka broker in minutes, assess fit with kmq assess kafka, then move topics and consumer-group offsets with kmq migrate.
Drop-in level: endpoint-only. KubeMQ speaks the native Kafka wire protocol — real
librdkafka/kcat/Java clients connect unchanged; you repoint bootstrap.servers, with no
client-library swap or code change. Moving an existing Apache Kafka, Amazon MSK, or Confluent
cluster onto KubeMQ is more than pointing at a new endpoint if you want consumers to resume
without reprocessing history — that historical topic-data and offset move is what the separate
kmq migrate tool below handles. This guide covers the read-only fitness check, the four-phase kmq migrate tool,
per-source auth, the large-history MirrorMaker 2 hybrid, rollback, and a post-cutover smoke
test. Start with the quick start below to see your client work first.
Overview
Migrating an existing Kafka workload onto KubeMQ has two stages: assess, then
migrate. kmq assess kafka scans your source cluster read-only and reports, per topic,
whether it's a straight repoint or needs a workaround. kmq migrate then does the actual
work — copying topic history and consumer-group offsets to KubeMQ in four phases (assess →
replicate → translate → cutover) — so consumers resume where they left off instead of
reprocessing everything from scratch.
If the source uses Kafka ACLs, rebuild the equivalent KubeMQ authorization rules before cutover. KubeMQ does not import Kafka ACLs, and Kafka ACL administration calls do not create KubeMQ policy. Follow Translating Kafka ACLs.
Data loss: dry-run before any production cutover
Know kmq migrate's scope before you rely on it. Its engine — byte-fidelity
replication, exact per-record offset translation, the cutover-completeness gate, and
oversized-record block-and-report — is proven on a real 3-node next cluster: a SIGKILL
mid-replication loses nothing, and a kill during cutover never double-seeds an offset.
Cluster-level consumer-resume is proven with a real franz-go client; Java/kcat/librdkafka
consumer-resume is proven single-node only — the full multi-client (Java / librdkafka /
franz-go) cluster consumer-resume run is what remains. The MSK and Confluent auth procedures
below are docs-grounded: validated against a local Apache Kafka, with cloud-specific details
taken from each vendor's own documentation.
Do a staged dry-run before any production cutover. Run kmq migrate assess and
kmq migrate translate against a copy of your data first and confirm the verdicts and the
offset-map preview look right — before you point a single production consumer at the
KubeMQ target.
Quick start: repoint one client
Try the switch on your own machine before you plan the migration.
The Docker tab points one client at the KubeMQ you started with Try KubeMQ, turning its connector on where needed and keeping your evaluation or license and your data. The Kubernetes tab works on a cluster you installed with Install on Kubernetes. Use it to try the switch, not for production; see Plans compared for staging and production licenses. If your current broker runs on this machine, stop it first, because two programs cannot publish the same port. The kmq command runs natively on Windows; the other commands are bash, so on Windows run them in WSL (Windows Subsystem for Linux).
You started KubeMQ with Try KubeMQ or Install with Docker. Both publish the Kafka port, 9092, so there is no connector to turn on.
Point your client at KubeMQ
bootstrap.servers=your-kafka-cluster.example.com:9092bootstrap.servers=localhost:9092Try KubeMQ publishes 9092 on localhost only, so this quick start serves clients on this machine. For clients elsewhere, see Configuration reference for CONNECTORS_KAFKA_ADVERTISED_HOST and the listener address.
Send and receive one message
Install kcat (brew install kcat on macOS, sudo apt-get install kcat on Debian or Ubuntu).
Produce one record, then read it back with a new consumer group. -o beginning lets the new group see the record you just produced, and -q hides the consumer's status lines.
echo "hello kubemq" | kcat -b localhost:9092 -t orders -Pkcat -q -b localhost:9092 -G orders-group -o beginning -c 1 ordersThe producer prints nothing. You should see:
hello kubemqNothing changed on the server, so there is nothing to undo.
The Kafka connector is on by default on Kubernetes too. A port-forward cannot serve Kafka clients, so reach the cluster as Reaching Kafka on Kubernetes describes, then run step 2 of the Docker tab against that address.
Before you move production traffic
- Read What Differs from Kafka and run the fit check in Assess fit before you commit.
- Plan how traffic moves and how you roll back: Rollback.
- Install KubeMQ for real: Choose your path.
Storage engine (zero-config)
Kafka on KubeMQ requires the next storage engine — but you don't set it manually. On a
fresh deployment, enabling the Kafka connector with the engine left unset auto-selects
next and logs a NOTICE; there is no separate "set store.engine=next first" step.
| Case | Behavior |
|---|---|
| Fresh deployment, Kafka enabled, engine unset | Auto-selects next. Zero-config — nothing to set before you enable Kafka. |
Existing legacy cluster, Kafka left at its default | The server starts without Kafka: port 9092 stays closed and an ERROR in the log names the port and the reason. Kafka cannot be retrofitted onto legacy data — there is no in-place engine migration — so stand up a new cluster on next. |
Existing legacy cluster, Kafka enabled explicitly | Fails closed at boot with a config error naming the conflicting store directory. Same remedy: a new cluster on next. |
Explicit STORE_ENGINE=legacy | Same as the two rows above: Kafka off for the run if it was only on by default, a boot refusal if you enabled it explicitly. |
Pinning STORE_ENGINE=next (or store.engine: next in config.yaml) always works and
always wins — it skips the auto-select probe entirely, which is the predictable choice for
IaC/GitOps. See Storage Engines
for the full engine model, mode isolation, and durability guarantees.
Assess fit — kmq assess kafka
Before touching anything, run the read-only assessor against your source cluster:
kmq assess kafka --bootstrap your-broker:9092 [--tls --sasl-mechanism scram-sha-256 --sasl-username <user> --sasl-password <pass>]kmq assess kafka never produces, commits an offset, or creates a topic — it's safe to run
against a production cluster. It reports, per topic, one of four verdicts, plus a single
overall migratable verdict:
| Verdict | Meaning |
|---|---|
| READY | No obstruction found — a straight repoint. |
| CAVEAT | Migratable, with a documented constraint (see below). |
| UNKNOWN | The topic's config couldn't be read (for example, a restricted cluster). Fail-safe — never assumed migratable. |
| BLOCKED | A hard blocker (see below) — this topic won't migrate through kmq migrate. |
Hard BLOCKED:
- Compacted topics (
cleanup.policy=compact) — compaction leaves non-contiguous offsets (compaction holes) that the exact offset-translation path can't follow. Move these via bulk data copy (the MirrorMaker 2 hybrid below) and rebuild consumer state, or keep them on the source. - More than 256 partitions on a topic — KubeMQ's per-topic partition cap. There's no auto-repartition (re-hashing would break partition-pinning and per-key ordering); repartition on the source first.
CAVEAT, not a block:
max.message.bytesabove the 1 MiB target cap. Config isn't data — your actual records may all be within the limit. A record that is over the cap doesn't fail assessment; it blocks and reports at replicate time instead (see Migrate below).
Not an assess blocker at all:
replication.factor. Not a fitness question: KubeMQ accepts any replication factor up to its node count, and every partition is replicated to every node regardless of the number asked for. A source topic at RF 3 therefore migrates as-is to a cluster of three or more nodes; a single-node target refuses RF above 1, as a one-broker Kafka would.
Active consumer groups on the source are tagged "requires exact offset translation at
cutover" — stop their consumers before you run kmq migrate cutover.
The verdicts above are the same rubric rendered by the
fitness matrix (T1–T4) — kmq assess kafka
just maps it onto your actual topics, configs, and consumer groups.
Migrate — kmq migrate
kmq migrate mirrors the shape of a real migration, one command per phase:
| Phase | Command | Touches the source | Touches the target | Reversible |
|---|---|---|---|---|
| assess | kmq migrate assess | read-only | — | Yes — read-only |
| replicate | kmq migrate replicate | read-only | writes records | Yes — target only |
| translate | kmq migrate translate | read-only | — | Yes — preview only |
| cutover | kmq migrate cutover | read-only | writes group offsets | Yes — re-runnable |
Strict acknowledgment is required
Before running kmq migrate replicate or kmq migrate cutover, set
Store.NextAckPolicy=strict on every KubeMQ target node, or set
STORE_NEXT_ACK_POLICY=strict, and restart every node. The shipped default is fast.
The migration commands refuse a target unless every node reports strict. A migration state
file built while any target used fast cannot be made safe by changing the setting afterward:
re-create the target and run replication again before cutover.
What the target must be
Replicate and a real cutover also refuse a target that is not a KubeMQ cluster managed by the Kubernetes operator: they bind the migration to the operator's deployment and data-volume identity, so a Docker server cannot be a target. Before replicating:
- Create every target topic first. Replicate refuses a topic it did not find with an explicit
creation identity. Use
kmq migrate prepare --input FILEto preview, thenkmq migrate prepare --apply PLAN_FILE --approve-hash HASHwith the hash the preview printed. - Run from a network that reaches every broker. The commands dial each broker the target advertises, and each broker's management API. On Kubernetes that is inside the cluster by default; see Running the migration on Kubernetes.
The source is only ever read. Every phase dials the source with a read-only client — none of them produce, commit an offset, join a consumer group, or auto-create a topic there. All writes (records, seeded offsets) land on the KubeMQ target only.
The assess phase is covered in detail above — see
Assess fit — kmq assess kafka.
Replicate — copy the history
kmq migrate replicate \
--bootstrap source-broker:9092 [source auth flags] \
--target-bootstrap kubemq-broker:9092 [target auth flags] \
--target-api http://kubemq-0:8080 --target-api http://kubemq-1:8080 --target-api http://kubemq-2:8080 \
--state ./migration.state--target-api is required: name the management API of every target node (never a
load balancer), one URL per Kafka broker the target advertises, using that broker's advertised
host name. Replicate refuses to start unless each node reports
Store.NextAckPolicy=strict, and it records what it saw so cutover can verify it again. When
the target requires sign-in (the default), add --target-context NAME, and save one signed-in kmq
context for each --target-api URL, with exactly that address: a saved credential is bound to one
endpoint.
Copies every source topic-partition into KubeMQ byte-for-byte — key, value, headers, and
the original CreateTime — partition-pinned (a record's source partition is preserved
exactly; no key re-hash). It produces with acks=all and, in the same pass, records an exact
per-record source-offset → target-offset map to --state.
- Resumable and crash-safe. Re-running with the same
--stateresumes from the last durably-copied offset per partition; already-acked records aren't re-copied. - Oversized records block — they never skip. A source record larger than the target's message-size cap halts replication at that offset with a report (topic/partition/offset/size). The partition watermark doesn't advance past it, so nothing is silently dropped. Resolve the record (or raise the target's max-message-bytes, if appropriate), then re-run.
- Restrict scope with
--topic(repeatable); the default is every non-internal source topic.
Translate — preview the offset map (no writes)
kmq migrate translate --bootstrap source-broker:9092 --target-bootstrap kubemq-broker:9092 --state ./migration.stateFor each source consumer group, shows the target resume offset each committed offset maps to,
and flags any group that can't be cut over cleanly. This is a read-only preview of exactly
what cutover would seed — the second half of the staged dry-run described in the
Overview.
Cutover — seed the offsets
# stop the source consumers first, then:
kmq migrate cutover --bootstrap source-broker:9092 --target-bootstrap kubemq-broker:9092 \
--target-api http://kubemq-0:8080 --target-api http://kubemq-1:8080 --target-api http://kubemq-2:8080 \
--state ./migration.stateAs with replicate, a real cutover needs one --target-api per target node (plus --target-context when the target requires sign-in); --dry-run does not.
Reads each source group's committed offsets, translates them through the offset map, and seeds them on the KubeMQ target via an empty-group offset commit — before any consumer joins. Point your consumers at KubeMQ and they resume at the correct position.
- Fail-closed completeness — zero tail loss on acked/committed offsets. A group is seeded only once replication has durably reached at least that group's committed offset on every partition. If a group consumed past what's been replicated, cutover refuses it with a clear reason, rather than seeding a position that would silently skip the un-replicated tail. Replicate further, then re-run.
- Offset-based only. Cutover resumes by offset, never by timestamp — a target
timestamp-seek would resolve KubeMQ's own ingestion time, not the preserved source
CreateTime, and would mis-position a consumer. - Idempotent. Re-running overwrites the seeded offsets safely.
--dry-runtranslates and checks completeness without writing. An active source or target group is refused; stop its consumers first. The--forceflag is deprecated and refused — there is no override.
Share groups (KIP-932)
Share groups are generally available on KubeMQ, so a workload using Java's
KafkaShareConsumer, the Kafka console tools or franz-go moves by repointing
bootstrap.servers. kmq assess kafka lists the source's share groups as a caveat finding.
kmq migrate cutover seeds classic group offsets only, so a share group needs these steps:
- Upgrade every node first. On a cluster moving from a release before share groups became
generally available, follow the
share-group upgrade rule
before you set start offsets or share consumers connect. While
CONNECTORS_KAFKA_SHARE_FEATURE_VERSIONis still1, the server refuses to create a group's first start offset withUNSUPPORTED_VERSION. - Choose where each share group starts. A new share group starts at the latest offset,
as in Kafka: records already on the topic are not delivered to it. To drain the replicated
history instead, set
CONNECTORS_KAFKA_SHARE_AUTO_OFFSET_RESET=earlieston the server for every new group, or setshare.auto.offset.reseton one group withkafka-configs.sh --entity-type groups. - Set exact start offsets before the first consumer joins. Stop the source share
consumers, then run
kmq kafka share-groups reset-offsets <group> --topic <t> --to-earliest|--to-latest|--to-offset N --executeagainst KubeMQ. It creates the group's start offset on every partition where the group has none, and the first consumer starts there. It is refused while a consumer is attached.reset-offsetsconnects to KubeMQ's Kafka listener, which on Kubernetes is reachable only inside the cluster by default. - Check the differences. A dead consumer's records come back only when their acquisition
lock expires (a 30-second lock by default),
ShareFetchanswers at once instead of waiting, and the in-flight limit counts stored batches rather than records — see the full list.
Running the migration on Kubernetes
On Kubernetes the KubeMQ Kafka listener is reachable only inside the cluster by default (see
Install on Kubernetes), so a workstation
cannot run the commands above against it. kmq can run each stage as an in-cluster Job with
kmq migrate job submit and kmq migrate job status (see kmq CLI).
That flow needs:
- a plan from
kmq migrate plan, which reads both clusters, so it too must be authored from a network that reaches every target broker; - a profiles ConfigMap, a state PersistentVolumeClaim, and restricted worker and cutover ServiceAccounts, created in the namespace beforehand, plus a credential Secret when the target requires sign-in;
- cutover as a sequence of staged Job operations (preflight, live check, fence, tail, seed and the rest), not one command.
Contact support to plan an in-cluster migration.
Repointing an application that logs in
An application that authenticated to Kafka carries its login in its client settings —
security.protocol=SASL_SSL or SASL_PLAINTEXT, sasl.mechanism, and a username and password or
a JAAS line. Changing bootstrap.servers alone keeps all of that. What happens next depends on
whether the KubeMQ server has Kafka users configured.
The server has Kafka users (Connectors.Kafka.Credentials, a users file, or the default user
— see Authentication). Keep the client's settings and
give it a KubeMQ user for the same mechanism. PLAIN and both SCRAM variants carry over directly;
Confluent Cloud API keys become an ordinary PLAIN user; MSK IAM has no KubeMQ equivalent, so those
clients move to SCRAM or mTLS. SASL_SSL needs the TLS listener on 9093; SASL_PLAINTEXT uses
9092.
The server has no Kafka users — the default for a fresh server — and the client still sends a
login. The server answers the login with UNSUPPORTED_SASL_MECHANISM, which the
client reports as an authentication failure, and the server logs, at most once a minute:
kafka client attempted SASL login but this server has no Kafka credentials configured, so the login was refused — remove the client's security.protocol / sasl.* settings, or configure Connectors.Kafka.CredentialsEither do what it says — drop security.protocol and every sasl.* setting from the client — or
configure users on the server.
Leftover TLS settings are the other common carry-over: a client with
security.protocol=SSL pointed at the plaintext port 9092 fails its handshake. Point it at
9093 with TLS configured on the server, or remove the TLS settings.
Per-source auth
The migration mechanics are identical across sources — only the auth flags differ, both
for reading the source (assess/replicate/translate/cutover) and for the application's
eventual switch to KubeMQ.
| Source | Reading the source | After cutover, on KubeMQ |
|---|---|---|
| Amazon MSK | MSK IAM — --sasl-mechanism aws-msk-iam with static keys | SASL/SCRAM or mTLS |
| Confluent Cloud | SASL/PLAIN over TLS — API key/secret as username/password | SASL/PLAIN-over-TLS (rotate the API-key credentials) |
| Self-managed Apache Kafka | PLAIN, SCRAM-SHA-256, or SCRAM-SHA-512 (+ TLS as required) | SASL/SCRAM, PLAIN, mTLS, or OAUTHBEARER |
MSK IAM is --sasl-mechanism aws-msk-iam with static keys — the AWS default credential
chain (profile / SSO / IMDS) is not used; supply the keys explicitly:
kmq migrate assess \
--bootstrap b-1.your-cluster.kafka.us-east-1.amazonaws.com:9098 \
--sasl-mechanism aws-msk-iam \
--aws-access-key <access-key> \
--aws-secret-key <secret-key> \
[--aws-session-token <token>]The application-side switch is config-only: IAM/SigV4 on the source becomes SASL/SCRAM or
mTLS on KubeMQ. Your code is unchanged beyond the bootstrap.servers value you were already
changing at migration.
Confluent Cloud is SASL/PLAIN over TLS — not SCRAM — with the API key/secret as the username/password:
kmq migrate assess \
--bootstrap pkc-xxxxx.us-east-1.aws.confluent.cloud:9092 \
--tls \
--sasl-mechanism plain \
--sasl-username <api-key> \
--sasl-password <api-secret>After cutover, rotate the Confluent API-key credentials to KubeMQ SASL/PLAIN-over-TLS. If you run Schema Registry, it keeps working unchanged against KubeMQ.
Self-managed Apache Kafka uses whichever SASL mechanism your cluster already runs:
kmq migrate assess \
--bootstrap kafka.internal:9093 \
--tls --tls-ca /path/to/ca.pem \
--sasl-mechanism scram-sha-256 \
--sasl-username <user> \
--sasl-password <pass>Kerberos / GSSAPI caveat. The bridge can consume a GSSAPI-secured source where the underlying client library supports it — but KubeMQ itself doesn't serve Kerberos. The migrated workload authenticates to KubeMQ with SASL/SCRAM, PLAIN, mTLS, or OAUTHBEARER instead.
Large histories — MirrorMaker 2 hybrid
kmq migrate is the offset-exact vehicle — tuned for correct consumer resume, not for
saturating a multi-terabyte historical backfill from a single process. For a very large
history, use a hybrid:
- MirrorMaker 2 for the bulk data copy. MM2 replicates records into KubeMQ
byte-perfectly, and its distributed workers scale the throughput. Use the shippable
config: keep every
*.replication.factorat or below the number of KubeMQ nodes (a larger one is refused, as on Kafka) and setoffset-syncs.topic.location=target(the default writes that topic to your source cluster instead). kmq migrate cutoverfor the offsets. Don't rely on MM2 for consumer-offset translation — its own checkpoint translation does not produce a correct zero-loss resume position. Use MM2 for the data, then seed offsets withkmq migrateagainst the same target.
MM2 (data) + kmq migrate (offsets) gives you MM2's throughput and kmq migrate's exact
consumer resume.
Rollback
Migration is non-destructive to the source — every phase only reads it. Rollback is:
- Point
bootstrap.serversback to the original source cluster. - Resume the source consumers.
Because the source's data and committed offsets are untouched throughout, it remains a complete, consistent fallback until you decommission it. The KubeMQ target can be discarded and the migration re-run from scratch.
OAUTHBEARER onboarding
If your Kafka clients already authenticate with OAUTHBEARER against an OIDC identity
provider, KubeMQ's Kafka connector accepts the same mechanism natively — there's no separate
credential model to bolt on for the migrated workload. See
the Kafka OAUTHBEARER settings for the
six-field OAuthBearer config block and its environment variables.
Verify the migration
For a first round trip, see the quick start; this section checks a finished migration.
Use a standard Kafka client pointed at the KubeMQ target to confirm the cutover worked.
-
Point a client at the KubeMQ target. Repoint
bootstrap.servers(orkcat -b) to the KubeMQ Kafka listener you migrated into. -
Produce a record to a migrated topic and confirm it's accepted:
kcat -b kubemq-broker:9092 -t orders -P <<< "smoke-test-record" -
Consume it back from the same topic:
kcat -b kubemq-broker:9092 -t orders -C -c 1 -
Check a consumer-group offset. Resume a migrated group and confirm it picks up at the translated offset — not
0:kafka-consumer-groups.sh --bootstrap-server kubemq-broker:9092 --describe --group my-groupThe committed offset shown should match
kmq migrate translate's preview for that group, and the group should resume without reprocessing already-consumed records.
See Also
Fitness matrix
What's drop-in, supported, roadmap, and unsupported when running Kafka workloads on KubeMQ (T1–T4).
Kafka
Adopt KubeMQ as a drop-in Kafka broker — what works, how to assess fit, and how to migrate.
Kafka configuration reference
Every Connectors.Kafka setting, env var, and CRD field — including OAUTHBEARER.
kmq CLI — Migration
The kmq assess kafka and kmq migrate command reference.
Was this page helpful?
TLS and mTLS
Secure the Kafka connector with TLS — the 9093 encrypted listener, server certificates, and mutual TLS where the certificate common name is the principal.
Reaching Kafka on Kubernetes
Why kubectl port-forward cannot work for Kafka clients, and the three ways that do: a client inside the cluster, one exposed node, or per-broker addresses.