KubeMQ
ConnectorsKafkaConcepts

Configuration

How the Kafka connector is enabled, ported, and secured — the on-by-default CONNECTORS_KAFKA_ENABLE flag, the 9092/9093 listeners, and the settings reference.

Overview

Configuring the Kafka connector is entirely a server-side decision. A Kafka client configures nothing KubeMQ-specific — it just points bootstrap.servers at the broker and, if the deployment requires it, supplies SASL credentials or a client certificate. Everything below — whether the listeners even open, which ports they bind, which authentication mechanisms are offered, and which storage engine backs the resulting topics — is decided once, on the server, and applies to every client that connects.

On by default

The connector ships enabled by default. A stock KubeMQ server binds port 9092 the moment it starts (9093 too, once TLS is configured) and advertises itself as a Kafka broker — there is no flag to set before a client can connect:

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

Setting CONNECTORS_KAFKA_ENABLE=false closes both listeners — a config-only rollback with nothing to migrate, because the connector never owned a separate data store to begin with: every produced record already lives in a plain Events Store log (see Architecture).

Because the connector is on by default rather than by explicit choice, it is never allowed to stop the server booting. If the configuration cannot run it — the store is pinned to the legacy storage engine, port 9092 or 9093 clashes with another enabled listener, a clustered deployment's raft members all sit on one host so per-broker addresses cannot be derived, or Broker.MaxPayload is too small — the server prints one stderr line and carries on without Kafka:

WARNING: the Kafka connector is on by default but cannot run with this configuration, so it is DISABLED for this run: ... — set Connectors.Kafka.Enable=false to silence this, or =true to make it a startup error

Setting CONNECTORS_KAFKA_ENABLE=true explicitly opts out of that tolerance: the same problems become fatal startup errors, which is what you want on a deployment that must not silently come up without its Kafka endpoint.

A single node with no HOST set advertises localhost:9092 to Kafka clients — fine for a local client, wrong for anything remote. Set CONNECTORS_KAFKA_ADVERTISED_HOST (and CONNECTORS_KAFKA_ADVERTISED_PORT) to the address remote clients can actually reach.

A multi-node deployment needs one operator habit up front: producers must use acks>=1. A single Kafka-facing Service can land a produce on any pod, and a follower forwards an acks>=1 produce to the leader transparently — but silently drops an acks=0 produce instead of forwarding it. Single-node deployments are unaffected.

Security posture

At a glance, the connector supports the same authentication and authorization shape a real Kafka deployment does, at the same layers:

  • SASL — PLAIN and both SCRAM-SHA-256/SCRAM-SHA-512 mechanisms are available once any Kafka credential is configured; a client authenticates with a username and password checked against the connector's own dedicated credential store, separate from KubeMQ's general-purpose auth.
  • OAUTHBEARER — OIDC-federated bearer tokens, validated against a configured issuer and offered only on the TLS listener — a bearer token is never accepted over plaintext.
  • mTLS — a client certificate presented on the TLS listener yields a principal derived from the certificate's common name, and only from a verified certificate chain.
  • ACLs — every request that reaches dispatch is authorized against KubeMQ's own policy engine, mapped onto the access level Kafka would expect: a produce or an offset commit needs write access, while a fetch or a group heartbeat needs read access.

None of this is mutually exclusive — a deployment can run SASL/SCRAM on the plaintext listener for internal traffic and OAUTHBEARER plus mTLS on the TLS listener for anything crossing a trust boundary. Authentication and TLS and mTLS walk through configuring each mechanism; Configuration reference has the copy-paste TOML/environment/Docker examples.

How configuration maps to behavior

Every Kafka setting can be supplied three ways — a [Connectors.Kafka] block in a TOML config file, a CONNECTORS_KAFKA_* environment variable, or (on Kubernetes) a typed field under spec.kafka on the KubeMQ cluster resource — and all three ultimately populate the same in-memory configuration the connector reads once at startup. That single source of truth is why the connector's runtime behavior is fully predictable from its configuration: the enable flag gates whether the listeners open at all, the port fields decide what a client dials, the credential and SASL-mechanism fields decide what the security posture above actually offers, and a handful of numeric fields — maximum connections, maximum message size, per-request fan-out caps — bound how much of the shared server the connector is allowed to consume.

The full field-by-field table — every setting, its default, its valid range, and its exact environment-variable and CRD names — lives in one canonical place: the Kafka settings reference. This page is orientation, not the source of truth for any individual field.

Kubernetes

On Kubernetes the connector is on by default as well — omitting the kafka block, or the enabled field, means on. The typed field is the way to turn it off:

# Helm values
kafka:
  enabled: false

The equivalent CRD field is spec.kafka.enabled: false on the KubemqCluster resource. On a fresh cluster nothing else is needed — no engine choice, no networking configuration — the connector opens an in-cluster, plaintext endpoint at <cluster>-kafka.<namespace>.svc:9092 that any in-cluster client can dial immediately. Reaching that endpoint from outside the cluster needs two more fields, advertisedHost and advertisedPort, plus a LoadBalancer or NodePort Service exposure — covered in Configuration reference and the Kafka settings reference, not restated here.

The next storage engine relationship

One configuration consequence deserves its own callout: enabling Kafka couples the deployment to KubeMQ's next storage engine, because Kafka's headline behaviors — compacted topics and the quorum-fsynced acknowledgment contract — exist only there. You do not configure this coupling directly. On a fresh store, Kafka auto-selects next automatically; on a store that already has data under the other engine, the default-on connector is skipped for that run with the stderr WARNING shown above, and an explicit CONNECTORS_KAFKA_ENABLE=true turns that into a clear configuration error instead — the connector never runs in a reduced mode. The full decision tree — fresh store, existing store, explicit pin — lives at Storage Engines → Zero-config engine selection.

Was this page helpful?

On this page