Kafka
Changelog (last updated v2.2)
- v2.2: single-writer support — two topics per database
-
Single-writer indexing requires two Kafka topics per database: a source log for client writes and a replica log for the indexing leader’s resolved output. XTDB uses Kafka’s consumer-group rebalance protocol to elect the leader for each database, and fences split-brain writes to the replica log by term — see ‘Leader election and fencing’ below.
Previously, a database used a single Kafka topic, and every indexer node consumed it independently. With single-writer, only the elected leader consumes the source topic; followers tail the replica topic instead.
Upgrading:
- The replica topic defaults to
${topic}-replicaand auto-creates whenautoCreateTopicis enabled, so existing deployments need no configuration changes to pick it up. - If multiple XTDB deployments share a Kafka cluster, give each a distinct
groupIdon its!Kafkacluster config — otherwise their consumer groups collide and leadership is assigned across deployment boundaries. - ACL-restricted topics need
Describe/Read/Writeon both the source and replica topics.
- The replica topic defaults to
- v2.2:
logClustersrenamed toremotes -
The Kafka cluster is now declared under
remotesrather thanlogClusters.logClustersis deprecated but still honoured, so existing config keeps working — rename toremoteswhen convenient. - v2.1: multi-database support
-
As part of multi-database support,
logClusterswere extracted in v2.1.Prior to that, the configuration in
logClusterswas within thelog:log: !KafkabootstrapServers: "localhost:9092"topic: "xtdb-log"# autoCreateTopic: true# pollDuration: "PT1S"# propertiesFile: "kafka.properties"# propertiesMap:# becamelogClusters:kafkaCluster: !KafkabootstrapServers: "localhost:9092"# pollDuration: "PT1S"# propertiesFile: "kafka.properties"# propertiesMap:log: !Kafkacluster: kafkaClustertopic: "xtdb-log"# autoCreateTopic: true
Apache Kafka can be used as XTDB’s message log. Each database uses two Kafka topics — a source log for client writes and a replica log for the indexing leader’s resolved output — plus Kafka’s consumer-group protocol to elect the leader for that database automatically. See ‘Database architecture’ for the concepts; this page covers how to set Kafka up to back them.
-
Add a dependency to the
com.xtdb/xtdb-kafkamodule in your dependency manager. -
On your Kafka cluster, XTDB requires two topics per database — a source log and a replica log:
- Both can be created manually and provided to the node config, or XTDB can create them automatically.
- If allowing XTDB to create the topics automatically, ensure that the connection properties supplied to the XTDB node have the appropriate permissions to create topics — XTDB will create each with the expected configuration values (single partition,
LogAppendTimetimestamps). Auto-created topics are unreplicated, so create them yourself for production.
-
Configure the topics and the broker — see Settings for which of these XTDB sets for you and which are yours.
-
XTDB should be configured to use the topics, and the Kafka cluster they’re hosted on. It should also be authorised to perform all of the necessary operations on both.
- For configuring the Kafka module to authenticate with the Kafka cluster, use the
propertiesFileorpropertiesMapconfiguration options to supply the necessary connection properties. See the example configuration below. - If the Kafka cluster is using ACLs, the XTDB node needs:
Describe/Read/Writeon both the source and replica topics.
- For configuring the Kafka module to authenticate with the Kafka cluster, use the
Settings
Section titled “Settings”Both topics — source and replica — take the same settings.
XTDB applies the topic settings it depends on only when it creates a topic itself (autoCreateTopic: true), so a topic you pre-create is entirely yours to configure.
The one setting it verifies on a topic that already exists is the partition count; the node refuses to start otherwise.
| Setting | Scope | Set by | Value |
|---|---|---|---|
| partition count | topic | XTDB on create, verified on an existing topic | Exactly 1. A single partition is what makes the log strictly ordered and lets leader election assign it to one consumer at a time. |
message.timestamp.type | topic | XTDB on create, not verified afterwards | LogAppendTime, so a record’s timestamp is when the broker appended it rather than when a producer sent it. Set it yourself on a pre-created topic. |
| replication factor | topic | XTDB creates with 1 | Pre-create the topic with 3 or more for production — auto-create is unreplicated. |
min.insync.replicas | topic | You | > 1, to make writes quorum-acknowledged. |
retention.ms | topic | You | Messages need not live on the log permanently. The default of 1 day suits most deployments; 1 week is a reasonable starting point where extra caution against data loss is wanted. |
max.message.bytes | topic | You | The 1MB default is fit for purpose unless your transactions are larger. |
cleanup.policy | topic | You | Leave at the default delete — XTDB never reads compacted messages. |
offsets.retention.minutes | broker, cluster-wide | You | Governs how long the leader-election consumer group survives with every XTDB node down, after which termEpoch has to be raised — see ‘Recreating the consumer group’. Seven days by default. This is not a per-topic setting, it applies to every consumer group on the cluster, and managed Kafka services often fix it. |
XTDB also sets its own producer and consumer properties — idempotent, acks=all writes, read_committed reads, auto.offset.reset=none, cooperative sticky assignment, and offset commits that keep the leader-election group alive.
propertiesMap and propertiesFile can override these, but they are chosen deliberately and overriding them can break leader election.
Configuration
Section titled “Configuration”To use the Kafka module, include the following in your node configuration:
## We first declare the Kafka cluster under `remotes`:
remotes: # You can define multiple Kafka clusters here, and refer to them by name in the log configuration. # Here we define a single Kafka cluster named "kafkaCluster". kafkaCluster: !Kafka # -- required
# A comma-separated list of host:port pairs to use for establishing the # initial connection to the Kafka cluster. # (Can be set as an !Env value) bootstrapServers: "localhost:9092"
# -- optional # The maximum time to block waiting for records to be returned by the Kafka consumer. # pollDuration: "PT1S"
# Path to a Java properties file containing Kafka connection properties, # supplied directly to the Kafka client. # (Can be set as an !Env value) # propertiesFile: "kafka.properties"
# A map of Kafka connection properties, supplied directly to the Kafka client. # propertiesMap:
# Consumer-group ID used for per-database leader election (v2.2+). # Defaults to "xtdb" — set a distinct value per deployment when multiple XTDB # deployments share a Kafka cluster, so their consumer groups don't collide. # groupId: "xtdb"
## For the database, we then create a log using the Kafka cluster we just defined:
log: !Kafka # -- required
# The name of the Kafka cluster to use for the source log. cluster: kafkaCluster
# Name of the Kafka topic to use for the source log. # (Can be set as an !Env value) topic: "xtdb-log"
# -- optional
# The name of the Kafka cluster to use for the replica log (v2.2+). # Defaults to the same cluster as the source log. # replicaCluster: kafkaCluster
# Name of the Kafka topic to use for the replica log (v2.2+). # Defaults to "${topic}-replica". # replicaTopic: "xtdb-log-replica"
# Whether or not to automatically create the topics, if they do not already exist. # Applies to both the source and replica topics. # autoCreateTopic: true
# Declares that the consumer group backing leader election has been recreated (v2.2+). # Raise it — never lower it — when that happens; see 'Recreating the consumer group' below. # termEpoch: 0SASL Authenticated Kafka Example
Section titled “SASL Authenticated Kafka Example”The following piece of node configuration demonstrates the following common use case:
- Cluster is secured with SASL - authentication is required from the module.
- Topic has already been created manually.
- Configuration values are being passed in as environment variables.
remotes: kafkaCluster: !Kafka bootstrapServers: !Env KAFKA_BOOTSTRAP_SERVERS propertiesMap: sasl.mechanism: PLAIN security.protocol: SASL_SSL sasl.jaas.config: !Env KAFKA_SASL_JAAS_CONFIG
log: !Kafka cluster: kafkaCluster topic: !Env XTDB_LOG_TOPIC autoCreateTopic: falseThe KAFKA_SASL_JAAS_CONFIG environment variable will likely contain a string similar to the following, and should be passed in as a secret value:
org.apache.kafka.common.security.plain.PlainLoginModule required username="username" password="password";Leader election and fencing
Section titled “Leader election and fencing”The ‘Database architecture’ page describes XTDB’s single-writer indexing model in terms of properties — exactly one leader per database, automatic failover, followers as hot standbys. This section describes how the Kafka module actually enforces those properties.
Leader election via consumer groups
Section titled “Leader election via consumer groups”Every XTDB node running the Kafka log subscribes to each database’s source topic via a shared Kafka consumer group (groupId, defaulting to "xtdb").
Kafka’s rebalance protocol then assigns each topic’s single partition to exactly one consumer across the group — that consumer is the leader for that database.
When a leader node stops (crash, network partition, long GC pause) or a new node joins the group, Kafka triggers a rebalance and the assignment moves.
Because all databases on a node share one consumer group via a single underlying consumer, Kafka’s CooperativeStickyAssignor distributes leaderships evenly across the cluster — e.g. with three nodes serving three databases, each node ends up leader for one database and follower for the other two.
Fencing by term
Section titled “Fencing by term”Each leader stamps every record it writes to the replica log with a term — the generation number Kafka assigns its consumer-group membership, which strictly increases each time the assignment moves. A leader applies and acknowledges a write only once it has read that write back from the replica log at its own term, with no higher-term record ahead of it.
If a rebalance has meanwhile moved leadership on, the incoming leader writes at a higher term. The outgoing leader reads that higher-term record back, recognises it has been superseded, and resigns — its own unconfirmed writes are never acknowledged. Followers apply the highest term they have seen and discard lower-term records. So at most one leader’s writes are ever confirmed for a given database, even across an unclean handover — without relying on Kafka transactions.
Recreating the consumer group (v2.2+)
Section titled “Recreating the consumer group (v2.2+)”A term is not that generation number alone: it pairs it with a termEpoch, which orders one incarnation of the consumer group against the next.
The generation orders elections within one incarnation; the epoch is what survives the group being recreated.
generationId restarts at 1 whenever the group is recreated, below the terms already on the replica log.
The broker deletes a consumer group once it has no members left and its committed offsets have expired, so how long a group survives with every XTDB node down is governed by the broker’s offsets.retention.minutes — seven days by default.
This turns on members, not traffic: a cluster that is up with no writes flowing keeps its group indefinitely, however long it idles.
offsets.retention.minutes is a broker setting rather than a topic one, so raising it for longer planned outages affects every consumer group on the cluster — and managed Kafka services often fix it.
Raise termEpoch on the log config whenever the group is recreated anyway — an outage past the retention period, a groupId change, or a group you delete deliberately:
log: !Kafka cluster: kafkaCluster topic: "xtdb-log" termEpoch: 1A node whose term is already fenced refuses to lead rather than indexing into a log every reader ignores, and its error names the current epoch:
leader term 0.1 is already fenced by 0.9 on the replica log — the leader-electioncounter has regressed (a recreated Kafka consumer group, or a restarted local log),so bump the log's termEpoch above 0Raise termEpoch, never lower it: a lower value puts the new leader back below the terms on the log.
Sharing a Kafka cluster across deployments
Section titled “Sharing a Kafka cluster across deployments”Leader election runs through a Kafka consumer group (groupId, defaulting to "xtdb").
If you run multiple XTDB deployments against the same Kafka cluster (e.g. staging + prod, or multiple tenants), give each a distinct groupId — otherwise they join the same group and Kafka assigns their topics’ partitions across both deployments’ nodes.
remotes: kafkaCluster: !Kafka bootstrapServers: "localhost:9092" groupId: "prod"Kafka Log Durability
Section titled “Kafka Log Durability”Kafka-backed logs offer strong durability, but require tuning and backup strategies to align with your recovery objectives.
Recommended Kafka Settings
Section titled “Recommended Kafka Settings”The replication factor, min.insync.replicas and retention.ms are the three that bear on data loss, and all three are yours rather than XTDB’s — see Settings.
Size retention.ms and retention.bytes so that unindexed messages survive long enough to be backed up or flushed.
See Apache Kafka documentation for details.
Managed services like Confluent Cloud may offer higher guarantees and simplified observability.
Strategies for Kafka Log Backup
Section titled “Strategies for Kafka Log Backup”There are three main ways to safeguard your XTDB Kafka log:
Point-in-Time Backups
Section titled “Point-in-Time Backups”Continuous Replication
Section titled “Continuous Replication”Use Kafka-native tools to replicate log data between clusters:
This allows for:
- Geo-redundancy
- Low-RPO disaster recovery
- Hot-standby clusters
Note: Replication does not replace backups --- it only increases availability.
Application-Level Transaction Replay
Section titled “Application-Level Transaction Replay”XTDB can rebuild its state from upstream sources (event logs, message queues) used to submit transactions.
Advantages:
- Independent recovery source
- Replay can be filtered, transformed, or validated
- Fills gaps between backup and failure