Authors: Luke Chen, Federico Valeri, Omnia Ibrahim, PoAn Yang, Kuan-Po Tseng, Jiunn-Yang Huang

Status

Current state: Under Discussion

Discussion thread: here

JIRA:

Please keep the discussion on the mailing list rather than commenting on the wiki (wiki discussions get unwieldy fast).

Motivation

Kafka deployments often require replicating data across geographically distributed clusters for disaster recovery (DR), regulatory compliance, data locality, cluster migrations or active-active architectures. While MirrorMaker 2 (MM2) provides cross-cluster replication capabilities, it presents significant operational challenges.

Goals

Cluster Mirroring addresses these operational challenges by integrating cross-cluster replication directly into Kafka brokers, providing a simpler and more robust solution for cross-cluster replication.

Figure 1: Cluster Mirroring Setup.

While Cluster Mirroring is optimized for geo-replication, DR and migration use cases where a single source cluster replicates to one or more destination clusters, its coordinator-based architecture provides a foundation for more complex topologies.

Non-Goals

Synchronous Replication

This proposal describes asynchronous replication between clusters. Support for synchronous replication is deferred to future work.

Producers write to the source cluster and receive acknowledgments based on the source cluster's replication requirements (e.g. acks=all ensures replication to all in-sync replicas within the source cluster). Data is then asynchronously replicated to destination clusters with no impact on producer latency or throughput.

This decision reflects the reality that cross-datacenter network latency makes synchronous replication impractical for some deployments. Requiring synchronous acknowledgment from a geographically distant cluster would introduce significant latency (typically 50-200ms for inter-region replication), making it unsuitable for latency-sensitive applications.

Implications for DR use cases:

Asynchronous replication should provide the right balance for DR use cases where availability and performance of the primary cluster must not be compromised by cross-datacenter latency. Applications requiring zero data loss across cluster failures can wait for the follow-up KIP that will extend this design to support synchronous mirroring, or handle the lag using application-level caching.

Stretched clusters are not suitable for DR scenarios because they provide no protection against software failures, configuration errors, and accidental topic deletions. Vendors that recommend stretched cluster deployments typically position them for high availability (HA) rather than DR, and notably, most do not offer stretched clusters as a managed service option, further underscoring the operational challenges and limited DR effectiveness of this architecture.

Unclean Leader Election

This proposal does not support unclean leader elections because there is no way to reconcile log divergence between source and destination clusters without a shared leader epoch. When the unclean.leader.election.enable is set to true, the broker will log a warning at every configuration synchronization period.

In normal Kafka operation, once a record is committed, part of the high watermark, it is immutable and will never be changed or removed. When a new leader is elected, followers use the epoch information to determine which records are safe to keep and which must be truncated to align with the new leader's log. Replicas eventually converge to the same data through epoch-based reconciliation. Unclean leader elections break this guarantee by allowing non-ISR brokers to become leaders, potentially with fewer records than were previously committed.

Source and destination clusters have completely independent controller architectures. Leadership changes in the source cluster happen independently of destination leadership changes. This means that epoch values diverge between clusters even though they represent the same logical topic partition. Source cluster epoch N and destination cluster epoch N have no inherent relationship, they represent different leadership events that happened at different times. This means that standard epoch comparison is insufficient because epochs are meaningful only within their originating cluster.

Solving this issue would require creating a shared leader epoch between source and destination clusters. Every time there is a source leader election we would need to notify the destination cluster and append data only after receiving a reply. This means that the overall latency would be cross-cluster replication latency plus intra-cluster replication latency. Read more in the Rejected Alternatives section.

Proposed Changes

Cluster Mirroring introduces a coordinator-based architecture integrated into Kafka brokers for managing cross-cluster replication. The design consists of three primary components that work together to provide automatic metadata synchronization and data replication. The following diagram illustrates how these components are wired together.

Figure 2: High Level Architecture.

The mirror name is stored as a topic-level internal configuration called mirror.name (same validation rules as topic name) that propagates through Kafka's metadata log as configuration change records. When topics are added to a mirror, the controller generates configuration metadata records that are replicated to all brokers through the standard metadata update mechanism.

Brokers monitor these configuration changes to detect when partitions they lead belong to a mirror, triggering the creation of mirror fetchers and enforcement of read-only semantics. This design ensures that mirror associations are visible, auditable, and manageable through standard Kafka configuration introspection tools while maintaining strict control over how mirroring relationships are established and modified.

Main Components

MirrorCoordinator

The MirrorCoordinator (MC) manages Cluster Mirroring state using a partitioned coordinator pattern similar to the group and transaction coordinators.

We use a composite key (mirror name, topic id, and partition number) to distribute coordination work across the __mirror_state topic's partitions, which is the internal compacted topic used to store mirror metadata. Each mirror partition independently hashes to a coordinator, spreading the load across all brokers in the cluster. This means a mirror with hundreds of partitions will have its state management distributed evenly rather than concentrated on a single broker.

Responsibilities:

Figure 3: Mirror Partition Lifecycle.

States descriptions:

Scenarios:

Starting a mirror (UNKNOWN -> PREPARING -> MIRRORING): The addTopicsToMirror command sets mirror.name config via the controller. The metadata update propagates to brokers. The broker leading the partition finds out the partition state via the coordinator, and this might trigger readMirrorState RPC to query from the remote coordinator and transitions to PREPARING if it's in a valid transition state (e.g. UNKNOWN). After truncation completes, it moves to MIRRORING and starts the mirror fetcher to fetch data from the source cluster.

Failover (MIRRORING -> STOPPING -> STOPPED): The removeTopicsFromMirror command appends the ".removed” suffix to the mirror.name config. The partition leader  detects the stop request, transitions to STOPPING, persists the last offset, then moves to STOPPED. The topic is now writable after the STOPPED state.

Restarting a stopped mirror (STOPPED -> PREPARING -> MIRRORING): The mirror.name config is set again. onMetadataUpdate sees the partition in STOPPED state and transitions to PREPARING, re-truncating and resuming replication.

MirrorMetadataManager

The MirrorMetadataManager (MMM) implements periodic metadata synchronization between source and destination clusters. It maintains persistent network connections to all source clusters.

Responsibilities:

  1. Connection Management: The manager maintains a connection pool with one blocking sender per source cluster. These connections are created lazily when the first topic for a mirror is added. Each sender uses the security credentials and network settings from the mirror configuration, allowing different mirrors to use different authentication mechanisms.
  2. Topic Metadata Synchronization: Every refresh cycle, the manager fetches topic metadata from source clusters using standard MetadataRequest calls. For each topic in the mirror configuration:
    1. Topic Creation: If a topic exists in the source but not the destination, the manager sends a CreateTopics request to the controller with identical partition count and configurations.
    2. Partition Expansion: If the source topic has more partitions than the destination, the manager sends a CreatePartitions request to scale up the destination topic to match.
    3. Configuration Sync: Topic configurations are compared between source and destination. Any differences trigger an IncrementalAlterConfigs request to align destination configs with the source.
    4. Topic Deletion: When a topic is deleted on the source cluster, the mirror partitions on the destination cluster moves to STOPPED state. This prevents accidental deletions to affect the destination cluster. In case it was intentional, the operator would need to manually remove the topic from the mirror.
  3. Consumer Group Offset Synchronization: The manager synchronizes classic and share consumer group offsets to enable seamless failover (no offset translation):
    1. Lists all consumer groups using ListGroups request.
    2. Fetches committed offsets for each group using OffsetFetch request or DescribeShareGroupOffsets request.
    3. Commits those offsets to the destination cluster’s group coordinator using the internal OffsetCommit or AlterShareGroupOffsets request.
  4. ACL Synchronization: Access control lists are mirrored from source to destination to maintain consistent security policies:
    1. Fetches all ACLs from the source using DescribeAcls request.
    2. Compares with the destination cluster’s current ACLs from the metadata image.
    3. Creates missing ACLs using CreateAcls request.
    4. Deletes ACLs that exist in destination but not in source using DeleteAcls request.

Cluster Mirroring allows users to modify configurations in the destination cluster, though these changes are periodically overridden by the topic configuration synchronization cycle. This design choice was made because while dynamic configuration changes could be blocked, static configuration changes via properties files cannot be prevented, making override inevitable. 

However, this approach presents challenges in environments with external governing systems like the Strimzi operator, where the continuous reconciliation process conflicts with the refresh cycle, potentially causing performance impacts. More critically, temporary configuration mismatches such as reduced retention periods or altered partition counts could lead to data loss or missing partitions until the next synchronization cycle detects and corrects the discrepancy, highlighting the need for careful operational awareness when mixing mirroring with external cluster management solutions.

Metadata synchronization operates at the mirror level rather than the partition level, so it uses a separate coordinator assignment based on the mirror name alone. Only the broker assigned as the metadata coordinator for a given mirror performs synchronization, and it applies changes only to the mirror partitions it manages. This avoids both redundant synchronization across brokers and unnecessary updates to partitions managed by other coordinators.

Each mirror can define its own filtering rules independently, loaded from the manager at each refresh cycle:

MirrorFetcherThread

The MirrorFetcherManager (MFM) extends AbstractFetcherManager to handle fetcher thread lifecycle for mirror partitions. It uses a three-dimensional key (fetcher ID, source broker endpoint, mirror name) to organize threads, ensuring that:

The MirrorFetcherThread (MFT) is a specialized implementation of AbstractFetcherThread that handles cross-cluster data replication with consumer Fetch requests and different epoch semantics than standard intra-cluster replication, but keeping the same log consistency validations. The destination cluster's replica is not registered as a follower in the source cluster. Using a follower Fetch request would cause the source broker to attempt updating follower replica status for a replica it doesn't know about. A consumer Fetch request avoids this issue, as it carries no such side effects on the source broker's replica state. In other words, destination partition leaders operate in a dual-role. They act as followers when fetching committed data from the source cluster leader up to the last stable offset (LSO), while simultaneously serving as leaders for their local replicas in the destination cluster. To maintain data consistency, destination partitions are read-only and reject produce requests from clients with ReadOnlyTopicException.

A mirror topic is created with the same topic ID as in the source cluster. This serves two purposes: it satisfies fetch request validation on the source broker, and it enables identity verification during failback where the destination cluster can confirm it is working with the exact same topic by comparing topic IDs.

A mirror leader partition begins fetching with an unknown source leader epoch. When it sends Fetch requests to the source cluster, the source leader may respond with a FencedLeaderEpochException. When such an error occurs, the mirror fetcher extracts the current source leader epoch from the error response, or from a separate metadata request if the source cluster doesn’t support fetch API v12, and updates its internal fetch state to track the source cluster's actual leader epoch. The last fetched epoch is always set to empty to disable log divergence checks due to unclean leader election (see non-goals section).

Figure 4: Mirror Leader Fetch State.

On subsequent Fetch requests:

  1. Fetch validation: The tracked source epoch ensures that fetched batches from the source cluster are validated against the correct source leader epoch, preventing acceptance of stale or invalid data (see KAFKA-18723).
  2. Epoch rewriting: When records are appended to the destination log, the batch epochs are rewritten to match the destination cluster's leader epochs, maintaining consistency within the destination cluster.

The source epoch tracking is purely for fetch validation, while the destination uses its own independent epoch sequence for replication and durability. This design keeps the two clusters' epoch spaces completely separate, allowing the destination to operate as a normal Kafka cluster with standard intra-cluster replication.

When the source partition's leader changes, a NotLeaderOrFollowerException is returned. At this point, the mirror fetcher thread queries the MirrorMetadataManager to get the updated endpoint and either creates a new fetcher thread or reuses one that is already connected to the new endpoint. This allows mirroring to continue seamlessly despite leadership changes in the source cluster.

When users remove a topic from the mirror, the partition will be removed from the fetcher thread, and any late fetch responses will be skipped because the partition is not registered anymore in the fetcher thread.

Failover Process 

Failover is initiated by calling the RemoveTopicsFromMirror API, which appends a ".removed" suffix to the mirror.name internal config. This transitions the mirror topics from read-only to writable state after the stopping process completes gracefully.

When producers reconnect to the destination cluster after failover, they obtain new producer IDs which are separate from previously mirrored IDs, so they begin writing with fresh sequence numbers starting from 0.

Consumers can reconnect to the destination cluster using the same group ID, resuming from the last synchronized offsets, minimizing data re-processing or gaps. The transition is transparent from the consumer's perspective and offset management continues normally through the destination's group coordinator.

# 9091 (source) -----> 9094 (destination)
# in case of disaster, the operator can failover by running the following command
bin/kafka-mirror.sh --bootstrap-server :9094 --remove --topic .* --mirror my-mirror
# 9091 (source) --x--> 9094 (destination)
# now all mirror topics are detached from the source cluster and accept writes (the two cluster are allowed to diverge)

Failback Process

Failback enables mirroring to be reversed after a failover, allowing the original source cluster to become the destination and vice versa. This is critical for scenarios where you want to fail back to the original cluster after recovering from an outage or planned maintenance.

For each partition, we track the high watermark (HW) by storing it in the cluster metadata as last mirrored offset (LMO) when removing a topic from a mirror (failover phase). The LMO represents the last record successfully mirrored from the original source cluster to the destination cluster before failover.

When failback is initiated on the old source cluster, it needs to determine where to truncate its log before starting to fetch from the new source cluster. If the new API is supported, the broker sends a LastMirrorredOffsets request to the new source cluster asking for the LMO, and then truncates its local log to the returned offset. If the new API is not supported, the broker truncates to zero and starts mirroring from scratch.

Before transitioning a mirror partition from PREPARING to MIRRORING, the MirrorCoordinator must ensure that all in-sync replicas (ISR) in the destination cluster have truncated their logs to the correct offset. If less than min ISR are available, we will skip and retry in the following fetch. This coordination step validates that every ISR member has completed truncation before the partition is allowed to begin actively fetching from the source cluster. Without it, the mirror leader could start appending new data from the source while local followers still hold divergent log segments, causing inconsistencies within the destination cluster. After truncation, reverse mirroring begins normally. Note that the log truncation on the reverse mirroring may cause the data loss for the records that didn’t get mirrored to the old destination cluster earlier.

# when the source cluster is back, the operator can failback by creating a mirror with the same name
echo "bootstrap.servers=localhost:9094" > /tmp/my-mirror.properties
bin/kafka-mirrors.sh --bootstrap-server :9091 --create --mirror my-mirror --mirror-config /tmp/my-mirror.properties
bin/kafka-mirrors.sh --bootstrap-server :"9091 --add --topic .* --mirror my-mirror
# 9091 (destination) <----- 9094 (source)

Existing Features Integration

Batch Compression

Cluster Mirroring preserves the compression format of record batches from the source cluster without recompression. When mirroring data, compressed record batches are copied directly from the source to the destination cluster, maintaining the original compression type (gzip, snappy, lz4, zstd, or none) and the exact byte-level representation of the data. This approach avoids unnecessary CPU overhead from decompression and recompression during replication, ensures bit-for-bit data integrity, and prevents potential issues with different compression implementations producing different outputs for the same data.

Topic Compaction

Cluster Mirroring fully supports log compacted topics, preserving both compacted records and offset gaps from the source cluster. When a topic uses cleanup.policy=compact, Kafka removes obsolete records with duplicate keys, creating gaps in the offset sequence. For example, if a source partition contains offsets 0-100 and compaction removes records at offsets 30-40 and 60-70, the remaining records will have gaps: offsets 0-29, 41-59, and 71-100 are missing.

The mirror leader replicates these compacted log segments exactly as they exist in the source cluster, maintaining the same offset assignments and gaps. After failover, when the mirror topic becomes writable, log compaction continues normally in the destination cluster according to the topic's compaction policy, and any new records produced locally will fill in after the highest mirrored offset.

When a destination cluster lags behind the source on a compacted topic, tombstone records may not yet be replicated at the time of failover. For example, if the source has a tombstone at offset 100 that deletes a key originally written at offset 3, but the destination has only replicated up to offset 50, that tombstone is never applied on the destination. After failover, offset 3 remains as a stale entry that will never be cleaned up. The problem is compounded on failback: truncating the former source to match the destination also discards the tombstone, so the orphaned key persists in both clusters permanently. This is an inherent limitation of asynchronous replication. Compaction correctness depends on the full sequence of tombstones being present, and any that fall beyond the replication watermark at failover time are lost. The same problem exists with MirrorMaker 2. The follow-up KIP for synchronous mirroring may address this by ensuring zero lag at switchover, guaranteeing all tombstones are replicated before failover occurs.

Topic Retention

Cluster Mirroring handles topic retention policies by periodically synchronizing the topic configurations from the source cluster, ensuring that the topic retention policies are consistent. When the source cluster applies retention policies, older log segments are deleted and the log start offset advances. For example, if a topic originally contained offsets 0-100 and retention deletes offsets 0-99, the source cluster's log start offset becomes 100. When the mirror leader fetches from the source, it discovers the new log start offset and updates its local log start offset to match, creating the same offset gap.

If a mirror follower attempts to fetch from an offset below the source cluster's log start offset (e.g. fetching offset 50 when log start offset is 100), the source broker returns an OffsetOutOfRangeException. The mirror leader handles this by truncating its local log to the source's current log start offset and resuming fetching from that point. This ensures the destination cluster mirrors the current retention state of the source cluster without attempting to replicate already-deleted data.

Consumer Groups

Cluster Mirroring synchronizes consumer group offsets from the source cluster to the destination cluster, enabling consumers to resume consumption from their last committed offset after failover. The MirrorMetadataManager periodically fetches consumer group committed offsets from the source cluster and replicates it to the destination cluster's. This ensures that consumer groups maintain their consumption progress across both clusters.

During offset synchronization, the committed offset in the destination cluster may temporarily exceed the current log end offset (LEO) of the mirror topic. For example, if a consumer commits offset 100 in the source cluster but the destination cluster has only mirrored up to offset 80 (LEO = 80), the MirrorMetadataManager still commits offset 100 to the destination cluster. This is acceptable because the mirror leader continues fetching data and the LEO will eventually advance to include offset 100. However, if a failover occurs before the mirrored data catches up, consumers attempting to resume from offset 100 will receive an OffsetOutOfRangeException. To handle this scenario gracefully, consumers should configure auto.offset.reset=latest when consuming from mirror topics. This ensures that if a committed offset is beyond the current LEO after failover, the consumer automatically resets to the latest available offset rather than failing or resetting to the earliest offset.

Security

Cluster Mirroring supports comprehensive security controls through both authorization and authentication mechanisms. On the destination cluster, mirror-related operations (creating mirrors, adding/removing topics from mirrors, managing mirror configurations) require the CLUSTER_ACTION permission on the cluster resource. This ensures that only authorized principals can establish and manage cluster mirrors. When configuring a mirror, operators specify ACLs that should be synchronized from the source cluster, and these ACLs are periodically replicated to the destination cluster to maintain consistent access control policies across both environments.

For connecting to the source cluster, Cluster Mirroring requires only the bootstrap server address and appropriate credentials, no other sensitive cluster information is exposed or required. The destination cluster's mirror configuration supports all standard Kafka authentication mechanisms including TLS/SSL for encrypted transport and SASL for client authentication.

Each mirror can be configured with its own security settings, allowing different mirrors to connect to source clusters with varying security requirements. This enables secure cross-cluster replication even when source and destination clusters use different authentication protocols or when connecting across security boundaries such as on-premises to cloud environments. All credentials are stored in the destination cluster's mirror configuration and used exclusively for establishing authenticated connections to the source cluster.

Idempotent Producer

The idempotent producers rely on producer IDs to detect duplicate writes and ensure idempotent production. To avoid conflicts with the destination cluster's producer ID space, we rewrite source producer IDs to occupy the unused negative space by applying the formula: 

destinationProducerId = -(sourceProducerId + 2)

The rationale of this formula is to keep the existing semantic of NO_PRODUCER_ID (-1) but still have a way to avoid the conflict. The CRC checksum is automatically recalculated after the producer ID changes to maintain batch integrity. Producer epochs from the source cluster are preserved exactly as they appear in the source batches. This ensures the last stable offset is correctly reflected because the producer state is updated after each append.

When a mirror topic becomes writable during failover, records with transformed producer IDs (<= -2) remain in the log with their original sequence numbers and epochs. Applications that reconnect to the destination cluster receive new producer IDs (>=0) from the destination's transaction coordinator, allowing them to continue producing.

Exactly-Once Semantics

Cluster Mirroring ensures transactional consistency when stopping by truncating to the LSO. Note that this doesn’t mean it supports exactly-once semantics (EOS) across clusters, which would require synchronous communication.

During the mirror stopping transition, the MirrorCoordinator performs a log truncation operation that resets each mirror partition to its LSO. This offset represents the point in the log where all transactions have been decided (committed or aborted), essentially the highest offset where data is known to be consistent from a transactional perspective. Any records beyond this point may belong to incomplete transactions and should not persist after mirroring stops. Note that the actual lag may be greater than what’s reported by the metrics.

This approach prevents a critical consistency issue: the destination cluster could retain partial transaction data that would never be completed since mirroring has stopped. This would leave the topic in an inconsistent state where read_committed consumers may be blocked due to incomplete transaction data. Additionally, the transaction coordinator would not be able to rollback these hanging transactions because there would be no __transaction_state metadata in the destination cluster.

Transactional Consumer Guarantees

Kafka consumers with isolation.level=read_committed determine transaction visibility using only the Last Stable Offset (LSO), which is computed from COMMIT/ABORT control markers in the log. Consumers never interact with the transaction coordinator or validate producer IDs. This separation between log-level markers (replicated) and coordinator state (not replicated) is why transactional consumers work correctly on mirror topics without mirroring coordinator state. The LSO truncation during failover ensures all remaining transactions have mirrored markers, maintaining this guarantee.

Example Scenario

Consider this source cluster log:

Offset

Type

isTxn

PID

Content

0

DATA_RECORD

true

4001

key=A, value=1

1

DATA_RECORD

true

4001

key=B, value=2

2

DATA_RECORD

true

4002

key=X, value=9

3

CONTROL_MARKER

true

4001

COMMIT marker for PID 4001

4

CONTROL_MARKER

true

4002

ABORT marker for PID 4002

5

DATA_RECORD

false

none

key=Z, value=10

If replication reaches offset 4 and the source cluster fails, the destination cluster contains data records for transaction 4002 (offset 2) without the abort marker (offset 4). This creates a hanging transaction that can never be committed or aborted on the destination cluster.

Note that this approach causes data loss for any in-flight transactions or non-mirrored completed transactions when we experiencing a lag during the failover and may result in already-processed records being lost if consumers on the destination cluster read uncommitted data.

Bandwidth Control

Cluster Mirroring adopts a dual-sided throttling mechanism that extends Kafka's existing bandwidth control capabilities to work across cluster boundaries.

  1. Destination Cluster Throttling: To avoid conflicts with intra-cluster replication controls, mirror-specific throttling configurations operate independently from standard replication throttling. The system provides two configuration levels: a broker-level rate limit (mirror.replication.throttled.rate) that sets the overall bandwidth ceiling for mirror replication traffic, and a topic-level replica list (mirror.replication.throttled.replicas) that specifies which partition-broker combinations should be throttled using the standard partition-index:broker-id notation. Operators can dynamically adjust throttling rates at runtime without restarting brokers, first setting a cluster-wide default rate, then fine-tuning specific topic partitions as mirroring progresses. This allows gradual bandwidth allocation as mirror relationships are established.
  2. Source Cluster Throttling: The source cluster side requires a different approach because mirror fetch requests operate as consumer traffic rather than replication traffic. This design is intentional since the mirroring must fetch only up to the LSO to maintain transactional consistency, which is a consumer-level guarantee not available through the replication protocol. Consequently, standard leader replication throttling mechanisms cannot apply to mirror traffic. Instead, the source cluster leverages Kafka's client quota system. Each mirror fetcher thread presents itself with a deterministic client identifier that encodes the broker ID, fetcher thread number, and mirror name. Operators can apply per-client byte rate quotas to these identifiers, effectively throttling the outbound mirror traffic from the source cluster. This approach integrates seamlessly with Kafka's existing quota enforcement infrastructure.

Tiered Storage

Tiered Storage is not initially supported, but a detailed design of the metadata synchronization protocol, API schema, and state management will be provided in a follow-up KIP. A mirror follower that receives an OffsetMovedToTieredStorageException from the source leader handles it by marking the partition as failed, and also the mirror partition state will move to FAILED state.

Share Group

Cluster Mirroring supports both traditional consumer groups and share consumer groups (Kafka Queue functionality) to ensure seamless failover for all consumer types. While the data mirroring mechanism remains identical, the offset synchronization strategy differs based on the group type.

Share consumer groups use a different offset management model based on Share-Partition Start Offset (SPSO) and Share-Partition End Offset (SPEO) rather than traditional committed offsets. First we retrieve the current SPSO for each share group using the DescribeShareGroupOffsets API from the source cluster, and then we update the SPSO in the destination cluster using the AlterShareGroupOffsets API, which also initializes the group state in both the group coordinator and share coordinator. This means the API can initialize a share group in the destination cluster even if it doesn't exist yet, eliminating the need for pre-creation or complex state management.

Kafka enforces that consumer group and share group names must be unique within a single cluster. This creates a potential conflict scenario during mirroring. When such conflicts occur, the offset commit operation will fail with GroupIdNotFoundException. Users must resolve these conflicts manually by either deleting the conflicting group in the destination cluster before mirroring begins, or excluding the conflicting groups from offset synchronization. These conflicts affect only offset synchronization and do not impact data mirroring itself. The topic data continues to replicate normally, and only the automatic offset synchronization for the conflicting groups is blocked.

Diskless Topics

At the time of writing, the Diskless Topics KIP (KIP-1500 and other sub-KIPs) are still under discussion, so there will be future KIPs to support this feature.

Active-Active Writes

Active-active topology is not initially supported in Cluster Mirroring, though it could potentially be achieved through topic prefixing and removing the reliance on topic ID for mirroring. This is a candidate for a future improvement KIP. 

Instead, bidirectional mirroring is supported, but only when mirroring different topics between clusters, allowing records produced to either cluster to be consumed from both. Unlike MirrorMaker 2, Cluster Mirroring does not need special cycle detection or prevention logic because the read-only enforcement inherently blocks the conditions that would create infinite replication loops.

Public Interfaces

Command-Line

A new dump flag allows to decode cluster mirroring metadata for debugging purpose:

$ bin/kafka-dump-log.sh --mirror-state-decoder --files /tmp/server*/data/__mirror_state-1/00000000000000000000.log
Dumping /home/fvaleri/Documents/kafka/build/test/server4/data/__mirror_state-0/00000000000000000000.log
Log starting offset: 0
baseOffset: 0 lastOffset: 0 count: 1 baseSequence: -1 lastSequence: -1 producerId: -1 producerEpoch: -1 partitionLeaderEpoch: 0 isTransactional: false isControl: false deleteHorizonMs: OptionalLong.empty position: 0 CreateTime: 1771239071191 size: 98 magic: 2 compresscodec: none crc: 3314855113 isvalid: true
| offset: 0 CreateTime: 1771239071191 keySize: 13 valueSize: 17 sequence: -1 headerKeys: [] key: {"type":"2","data":{"mirrorName":"my-mirror"}} payload: {"version":"0","data":{"topicName":"my-topic","partition":0,"state":0}}
baseOffset: 1 lastOffset: 1 count: 1 baseSequence: -1 lastSequence: -1 producerId: -1 producerEpoch: -1 partitionLeaderEpoch: 0 isTransactional: false isControl: false deleteHorizonMs: OptionalLong.empty position: 98 CreateTime: 1771239071219 size: 98 magic: 2 compresscodec: none crc: 3968746657 isvalid: true
| offset: 1 CreateTime: 1771239071219 keySize: 13 valueSize: 17 sequence: -1 headerKeys: [] key: {"type":"2","data":{"mirrorName":"my-mirror"}} payload: {"version":"0","data":{"topicName":"my-topic","partition":0,"state":1}}

A new command-line tool kafka-mirrors.sh provides administrative operations for managing cluster mirrors.

Create a new cluster mirror configuration in the destination cluster:

$ echo "bootstrap.servers=localhost:9092" >/tmp/mirror.properties
$ bin/kafka-mirror.sh --bootstrap-server :9094 --create --mirror my-mirror --mirror-config /tmp/mirror.properties
Created mirror my-mirror

Add a topic or set of topics to an existing cluster mirror (start mirroring):

$ bin/kafka-mirror.sh --bootstrap-server :9094 --add --topic my-topic --mirror my-mirror --replication-factor 2
Added 1 topic(s) to mirror my-mirror: [my-topic]

List configured mirrors with additional information:

$ bin/kafka-mirrors.sh --bootstrap-server :9094 --list
MIRROR                         TOPICS     SOURCE-BOOTSTRAP                                                                                                                 
my-mirror                      2          localhost:9091                                                                                                                   
new-mirror                     1          localhost:9091

Describe configured mirrors to check their lag compared to their source topics:

$ bin/kafka-mirrors.sh --bootstrap-server :9094 --describe
MIRROR                         TOPIC                                    PARTITION  SOURCE-OFFSET   DESTINATION-OFFSET LAG      STATE       
my-mirror                      bar                                      0          2324            2324               0        MIRRORING   
my-mirror                      foo                                      0          69              66                 3        MIRRORING   
my-mirror                      foo                                      1          94              84                 10       MIRRORING   
my-mirror                      foo                                      2          94              90                 4        MIRRORING   
new-mirror                     baz                                      0          189             189                0        MIRRORING   
new-mirror                     baz                                      1          859             859                0        MIRRORING

Pause/resume mirroring for a specific topic or set of topics (topics reamin read-only):

TODO

Remove a specific topic or set of topics from a mirror (failover: all partitions become writable):

$ bin/kafka-mirror.sh --bootstrap-server :9094 --remove --topic my-topic --mirror my-mirror
Removed 1 topic(s) from mirror my-mirror: [my-topic]

Delete a mirror including its topics and configuration (the mirror must be empty or include only stopped partitions):

TODO

Alter mirror configuration (e.g. authentication):

TODO

Throttling on the destination cluster:

$ bin/kafka-configs.sh --bootstrap-server :9094 --entity-type brokers --entity-name 4 \
  --alter --add-config mirror.replication.throttled.rate=100000000
Completed updating config for broker 4.

$ bin/kafka-configs.sh --bootstrap-server :9094 --entity-type topics --entity-name my-topic \
  --alter --add-config mirror.replication.throttled.replicas=[0:4]
Completed updating config for topic my-topic.

Throttling on the source cluster:

$ bin/kafka-configs.sh --bootstrap-server :9091 --alter --add-config 'consumer_byte_rate=1024' \
  --entity-type clients --entity-name broker-4-fetcher-0-mirror-my-mirror
Completed updating config for client broker-4-fetcher-0-mirror-my-mirror.

Admin Client

New methods are added to the Admin interface for programmatic cluster mirror management, along with their supporting classes:

CreateMirrorResult createMirror(String mirrorName, Map<String, String> configs, CreateMirrorOptions options);

AddTopicsToMirrorResult addTopicsToMirror(Map<String, String> topicToMirrorName, AddTopicsToMirrorOptions options);

RemoveTopicsFromMirrorResult removeTopicsFromMirror(String mirrorName, Set<String> topics, RemoveTopicsFromMirrorOptions options);

ListMirrorsResult listMirrors(ListMirrorsOptions options);

DescribeMirrorsResult describeMirrors(Collection<String> mirrorNames, DescribeMirrorsOptions options);

Protocol Changes

This KIP extends CreateTopic API, but also introduces some new APIs and metadata records.

CreateTopic

The CreateTopic API is extended to add information required for mirror topic creation.

// new added
{ "name": "MirrorInfo", "type": "MirrorInfo", "versions": "8+", "nullableVersions": "8+", "ignorable": true,
  "about": "Mirror information for creating a mirror topic from a source cluster.", "fields": [
    { "name": "TopicId", "type": "uuid", "versions": "8+",
    "about": "The topic ID from the source cluster." }
]}

The topic ID field ensures mirror topics retain the same topic ID as the source cluster topic. This allows fetch requests to pass validation on the source broker, and enables the system to verify that a topic being mirrored to a same-named topic in the destination cluster is indeed the same logical topic, not a name collision.

In normal topic creation, the MirrorInfo field will be null. When receiving the CreateTopic request, the controller will check the new field. If it is not set, the topic ID will be generated with random UUID as usual. Otherwise, the controller will do the following validation:

  1. This topic ID is not used by other topics in the current cluster
  2. The replicas for the partition assignment are all active and not in fenced or controlled shutdown. This is to make sure when a topic gets deleted and re-created with the same topic ID, the stale offline log dir won’t be treated as the active log dir after it becomes online (KAFKA-16234).

EntityType

A new entity type is added for  in the message generator to provide schema-level type validation for mirror name fields:

public enum EntityType {
    // ... existing types ...
    @JsonProperty("mirrorName")
    MIRROR_NAME(FieldType.StringFieldType.INSTANCE); // New type
}

CreateMirror

The CreateMirror API allows users to create a mirror and supply its configuration. When the broker receives the request, it validates that the mirror name is not already in use, contains only permitted characters, and does not end with the ".removed" suffix. Once validated, the request is forwarded to the controller, which persists the configuration in the metadata log.

CreateMirrorRequest

{
  "apiKey": TBD,
  "type": "request",
  "listeners": ["broker", "controller"],
  "name": "CreateMirrorRequest",
  "latestVersionUnstable": true,
  // Version 0 is the initial version.
  "validVersions": "0",
  "flexibleVersions": "0+",
  "fields": [
    { "name": "MirrorName", "type": "string", "versions": "0+", "entityType": "mirrorName",
      "about": "The cluster mirror name."},
    { "name": "Config", "type": "[]MirrorConfig", "versions": "0+",
      "about": "The cluster mirror configurations.",  "fields": [
      { "name": "Name", "type": "string", "versions": "0+", "mapKey": true,
        "about": "The configuration key name." },
      { "name": "Value", "type": "string", "versions": "0+", "nullableVersions": "0+",
        "about": "The value to set for the configuration key."}
    ]}
  ]
}

CreateMirrorResponse

{
  "apiKey": TBD,
  "type": "response",
  "name": "CreateMirrorResponse",
  // Version 0 is the initial version.
  "validVersions": "0",
  "flexibleVersions": "0+",
  "fields": [
    { "name": "ThrottleTimeMs", "type": "int32", "versions": "0+",
      "about": "The duration in milliseconds for which the request was throttled due to a quota violation, or zero if the request did not violate any quota." },
    { "name": "ErrorCode", "type": "int16", "versions": "0+",
      "about": "The error code, or 0 if there was no error." },
    { "name": "ErrorMessage", "type": "string", "versions": "0+", "nullableVersions": "0+", "ignorable": true,
      "about": "The error message, or null if there was no error." }
  ]
}

AddTopicsToMirror

The AddTopicsToMirror API adds topics to a specified mirror. The broker validates that all target topic partitions are in either UNKNOWN or STOPPED state; otherwise, the request is rejected with an INVALID_REQUEST error. Once validated, the request is forwarded to the controller, which sets the mirror.name topic config to the specified mirror name.

AddTopicsToMirrorRequest

{
  "apiKey": TBD,
  "type": "request",
  "listeners": ["broker", "controller"],
  "name": "AddTopicsToMirrorRequest",
  // Version 0 is the initial version.
  "validVersions": "0",
  "flexibleVersions": "0+",
  "fields": [
    { "name": "Topics", "type": "[]TopicData", "versions": "0+", "about": "The data for the topics.",
      "fields": [
        { "name": "TopicId", "type": "uuid", "versions": "0+", "about": "The unique topic ID."},
        { "name": "TopicName", "type": "string", "versions": "0+", "mapKey": true, "entityType": "topicName",
          "about": "The topic name." },
        { "name": "MirrorName", "type": "string", "versions": "0+", "entityType": "mirrorName",
          "about": "The cluster mirror name."}
      ]}
  ]
}

AddTopicsToMirrorResponse

{
  "apiKey": TBD,
  "type": "response",
  "name": "AddTopicsToMirrorResponse",
  // Version 0 is the initial version.
  "validVersions": "0",
  "flexibleVersions": "0+",
  "fields": [
    { "name": "ThrottleTimeMs", "type": "int32", "versions": "0+",
      "about": "The duration in milliseconds for which the request was throttled due to a quota violation, or zero if the request did not violate any quota." },
    { "name": "ErrorCode", "type": "int16", "versions": "0+",
      "about": "The error code, or 0 if there was no error." },
    { "name": "Topics", "type": "[]TopicResult", "versions": "0",
      "about": "The results for the topics.", "fields": [
      { "name": "Name", "type": "string", "versions": "0", "entityType": "topicName",
        "about": "The topic name." },
      { "name": "ErrorCode", "type": "int16", "versions": "0",
        "about": "The error code, or 0 if there was no error." }
    ]}
  ]
}

RemoveTopicsFromMirror

The RemoveTopicsFromMirror API allows users to detach topics from their associated mirror. The broker validates that all target topic partitions are in either PREPARING or MIRRORING state. Once validated, the request is forwarded to the controller, which appends the ".removed" suffix to the mirror.name topic config to mark the topics as no longer mirrored.

RemoveTopicsFromMirrorRequest

{
  "apiKey": TBD,
  "type": "request",
  "listeners": ["broker", "controller"],
  "name": "RemoveTopicsFromMirrorRequest",
  // Version 0 is the initial version.
  "validVersions": "0",
  "flexibleVersions": "0+",
  "fields": [
    { "name": "Topics", "type": "[]TopicData", "versions": "0+", "about": "The data for the topics.",
      "fields": [
        { "name": "TopicId", "type": "uuid", "versions": "0+", "about": "The unique topic ID."},
        { "name": "TopicName", "type": "string", "versions": "0+", "mapKey": true, "entityType": "topicName",
          "about": "The topic name." }
      ]}
  ]
}

RemoveTopicsFromMirrorResponse

{
  "apiKey": TBD,
  "type": "response",
  "name": "RemoveTopicsFromMirrorResponse",
  // Version 0 is the initial version.
  "validVersions": "0",
  "flexibleVersions": "0+",
  "fields": [
    { "name": "ThrottleTimeMs", "type": "int32", "versions": "0+",
      "about": "The duration in milliseconds for which the request was throttled due to a quota violation, or zero if the request did not violate any quota." },
    { "name": "ErrorCode", "type": "int16", "versions": "0+",
      "about": "The error code, or 0 if there was no error." },
    { "name": "Topics", "type": "[]TopicResult", "versions": "0",
      "about": "The results for the topics.", "fields": [
      { "name": "Name", "type": "string", "versions": "0", "entityType": "topicName",
        "about": "The topic name." },
      { "name": "ErrorCode", "type": "int16", "versions": "0",
        "about": "The error code, or 0 if there was no error." }
    ]}
  ]
}

LastMirroredOffset

The LastMirroredOffset API allows destination cluster partition leaders in PREPARING state to query the LMO from the source cluster. If the source cluster has no record of this offset in its internal topic, it returns 0, meaning the log must be truncated to the beginning and mirroring starts from scratch. This is particularly important during failback. The last mirrored offset identifies where mirrored data ends and un-mirrored data begins. Records beyond this offset must be truncated before mirroring new data from the new source cluster; otherwise, the two clusters would contain inconsistent data.

ListMirroredOffsetsRequest

{
  "apiKey": TBD,
  "type": "request",
  "listeners": ["broker", "controller"],
  "name": "LastMirroredOffsetsRequest",
  // Version 0 is the initial version.
  "validVersions": "0",
  "flexibleVersions": "0+",
  "fields": [
    { "name": "MirrorName", "type": "string", "versions": "0+", "entityType": "mirrorName",
      "about": "The cluster mirror name." },
    { "name": "Topics", "type": "[]TopicData", "versions": "0",
      "about": "The data for the topics.", "fields": [
      { "name": "Name", "type": "string", "versions": "0", "entityType": "topicName",
        "about": "The topic name." },
      { "name": "Partitions", "type": "[]PartitionData", "versions": "0",
        "about": "The data for the partitions.", "fields": [
        { "name": "PartitionIndex", "type": "int32", "versions": "0",
          "about": "The partition index." }
      ]}
    ]}
  ]
}

LastMirroredOffsetsResponse

{
  "apiKey": TBD,
  "type": "response",
  "name": "LastMirroredOffsetsResponse",
  // Version 0 is the initial version.
  "validVersions": "0",
  "flexibleVersions": "0+",
  "fields": [
    { "name": "ThrottleTimeMs", "type": "int32", "versions": "0+",
      "about": "The duration in milliseconds for which the request was throttled due to a quota violation, or zero if the request did not violate any quota." },
    { "name": "ErrorCode", "type": "int16", "versions": "0+",
      "about": "The error code, or 0 if there was no error." },
    { "name": "Topics", "type": "[]TopicResult", "versions": "0",
      "about": "The results for the topics.", "fields": [
      { "name": "Name", "type": "string", "versions": "0", "entityType": "topicName",
        "about": "The topic name." },
      { "name": "Partitions", "type": "[]PartitionResult", "versions": "0",
        "about": "The results for the partitions.", "fields": [
        { "name": "PartitionIndex", "type": "int32", "versions": "0",
          "about": "The partition index." },
        { "name": "LastMirroredOffset", "type": "int64", "versions": "0",
          "about": "The last mirrored offset." },
        { "name": "ErrorCode", "type": "int16", "versions": "0",
          "about": "The error code, or 0 if there was no error." }
      ]}
    ]}
  ]
}

ListMirrors

The ListMirrors API returns the current mirror names and their associated topic counts in the cluster.

ListMirrorsRequest

{
  "apiKey": TBD,
  "type": "request",
  "listeners": ["broker"],
  "name": "ListMirrorsRequest",
  // Version 0 is the initial version.
  "validVersions": "0",
  "flexibleVersions": "0+",
  "fields": []
}

ListMirrorsResponse

{
  "apiKey": TBD,
  "type": "response",
  "name": "ListMirrorsResponse",
  // Version 0 is the initial version.
  "validVersions": "0",
  "flexibleVersions": "0+",
  "fields": [
    { "name": "ThrottleTimeMs", "type": "int32", "versions": "0+", "ignorable": true,
      "about": "The duration in milliseconds for which the request was throttled due to a quota violation, or zero if the request did not violate any quota." },
    { "name": "ErrorCode", "type": "int16", "versions": "0+",
      "about": "The error code, or 0 if there was no error." },
    { "name": "Mirrors", "type": "[]ListedMirror", "versions": "0+",
      "about": "Each mirror in the response.", "fields": [
      { "name": "MirrorName", "type": "string", "versions": "0+", "entityType": "mirrorName",
        "about": "The cluster mirror name." },
      { "name": "SourceBootstrap", "type": "string", "versions": "0+",
        "about": "The source cluster bootstrap servers." },
      { "name": "TopicCount", "type": "int32", "versions": "0+", "default": "0",
        "about": "The number of topics configured for this mirror. 0 indicates an empty mirror with no topics." }
    ]}
  ]
}

DescribeMirrors

The DescribeMirrors API is to retrieve the information about the mirror names, including the partition state, lag, source offset and destination offset.

DescribeMirrorsRequest

{
  "apiKey": TBD,
  "type": "request",
  "listeners": ["broker"],
  "name": "DescribeMirrorsRequest",
  // Version 0 is the initial version.
  "validVersions": "0",
  "flexibleVersions": "0+",
  "fields": [
    { "name": "MirrorNames", "type": "[]string", "versions": "0+", "entityType": "mirrorName",
      "about": "The names of the mirrors to describe. Null or empty array means all mirrors." },
    { "name": "IncludeAuthorizedOperations", "type": "bool", "versions": "0+", "default": "false",
      "about": "Whether to include authorized operations." }
  ]
}

DescribeMirrorsResponse

{
  "apiKey": TBD,
  "type": "response",
  "name": "DescribeMirrorsResponse",
  // Version 0 is the initial version.
  "validVersions": "0",
  "flexibleVersions": "0+",
  "fields": [
    { "name": "ThrottleTimeMs", "type": "int32", "versions": "0+",
      "about": "The duration in milliseconds for which the request was throttled due to a quota violation, or zero if the request did not violate any quota." },
    { "name": "Mirrors", "type": "[]DescribedMirror", "versions": "0+",
      "about": "Each described mirror.", "fields": [
      { "name": "ErrorCode", "type": "int16", "versions": "0+",
        "about": "The error code, or 0 if there was no error." },
      { "name": "MirrorName", "type": "string", "versions": "0+", "entityType": "mirrorName",
        "about": "The cluster mirror name." },
      { "name": "Topics", "type": "[]TopicPartitions", "versions": "0+",
        "about": "Each topic in the mirror.", "fields": [
        { "name": "TopicName", "type": "string", "versions": "0+",
          "about": "The topic name." },
        { "name": "Partitions", "type": "[]PartitionDetail", "versions": "0+",
          "about": "Each partition detail.", "fields": [
          { "name": "PartitionIndex", "type": "int32", "versions": "0+",
            "about": "The partition index." },
          { "name": "SourceOffset", "type": "int64", "versions": "0+",
            "about": "The high watermark offset from the source cluster leader." },
          { "name": "DestinationOffset", "type": "int64", "versions": "0+",
            "about": "The log end offset on the destination cluster." },
          { "name": "Lag", "type": "int64", "versions": "0+",
            "about": "The lag (source offset - destination offset)." },
          { "name": "State", "type": "string", "versions": "0+",
            "about": "The partition state (INITIALIZING, PREPARING, MIRRORING, STOPPING, STOPPED, FAILED)." }
        ]}
      ]},
      { "name": "AuthorizedOperations", "type": "int32", "versions": "0+", "default": "-2147483648",
        "about": "32-bit bitfield to represent authorized operations for this mirror." }
    ]}
  ]
}

ReadMirrorStates

The ReadMirrorStates RPC reads mirror states from the coordinator broker when it resides on a different node than the requesting broker.

ReadMirrorStatesRequest

{
  "apiKey": TBD,
  "type": "request",
  "listeners": ["broker", "controller"],
  "name": "ReadMirrorStatesRequest",
  // Version 0 is the initial version.
  "validVersions": "0",
  "flexibleVersions": "0+",
  "fields": [
    { "name": "MirrorName", "type": "string", "versions": "0+", "entityType": "mirrorName",
      "about": "The cluster mirror name." },
    { "name": "Topics", "type": "[]TopicData", "versions": "0",
      "about": "The data for the topics.", "fields": [
      { "name": "Name", "type": "string", "versions": "0", "entityType": "topicName",
        "about": "The topic name." },
      { "name": "Partitions", "type": "[]PartitionData", "versions": "0",
        "about": "The data for the partitions.", "fields": [
        { "name": "PartitionIndex", "type": "int32", "versions": "0",
          "about": "The partition index." }
        ]}
      ]}
  ]
}

ReadMirrorStatesResponse

{
  "apiKey": TBD,
  "type": "response",
  "name": "ReadMirrorStatesResponse",
  // Version 0 is the initial version.
  "validVersions": "0",
  "flexibleVersions": "0+",
  "fields": [
    { "name": "ThrottleTimeMs", "type": "int32", "versions": "0+",
      "about": "The duration in milliseconds for which the request was throttled due to a quota violation, or zero if the request did not violate any quota." },
    { "name": "ErrorCode", "type": "int16", "versions": "0+",
      "about": "The error code, or 0 if there was no error." },
    { "name": "Topics", "type": "[]TopicResult", "versions": "0",
      "about": "The read results for the topics.", "fields": [
      { "name": "Name", "type": "string", "versions": "0", "entityType": "topicName",
        "about": "The topic name." },
      { "name": "Partitions", "type": "[]PartitionResult", "versions": "0",
        "about": "The results for the partitions.", "fields": [
        { "name": "PartitionIndex", "type": "int32", "versions": "0",
          "about": "The partition index." },
        { "name": "LastMirroredOffset", "type": "int64", "versions": "0",
          "about": "The last mirrored offset." },
        { "name": "State", "type": "int8", "versions": "0+",
          "about": "The mirror partition state." },
        { "name": "ErrorCode", "type": "int16", "versions": "0",
          "about": "The error code, or 0 if there was no error." }
      ]}
    ]}
  ]
}

WriteMirrorStates

The WriteMirrorStates RPC writes mirror state updates to the coordinator broker when it resides on a different node than the requesting broker.

WriteMirrorStatesRequest

{
  "apiKey": TBD,
  "type": "request",
  "listeners": ["broker", "controller"],
  "name": "WriteMirrorStatesRequest",
  // Version 0 is the initial version.
  "validVersions": "0",
  "flexibleVersions": "0+",
  "fields": [
    { "name": "MirrorName", "type": "string", "versions": "0+", "entityType": "mirrorName",
      "about": "The mirror name." },
    { "name": "Topics", "type": "[]TopicData", "versions": "0",
      "about": "The data for the topics.", "fields": [
      { "name": "Name", "type": "string", "versions": "0", "entityType": "topicName",
        "about": "The topic name." },
      { "name": "Partitions", "type": "[]PartitionData", "versions": "0",
        "about": "The data for the partitions.", "fields": [
        { "name": "PartitionIndex", "type": "int32", "versions": "0",
          "about": "The partition index." },
        { "name": "LastMirroredOffset", "type": "int64", "versions": "0",
          "about": "The last mirrored offset." },
        { "name": "State", "type": "int8", "versions": "0+",
          "about": "The mirror partition state." }
      ]}
    ]},
    { "name": "RemovedTopics", "type": "[]string", "versions": "0+", "about": "The topic names to be removed." }
  ]
}

WriteMirrorStatesResponse

{
  "apiKey": TBD,
  "type": "response",
  "name": "WriteMirrorStatesResponse",
  // Version 0 is the initial version.
  "validVersions": "0",
  "flexibleVersions": "0+",
  "fields": [
    { "name": "ThrottleTimeMs", "type": "int32", "versions": "0+",
      "about": "The duration in milliseconds for which the request was throttled due to a quota violation, or zero if the request did not violate any quota." },
    { "name": "ErrorCode", "type": "int16", "versions": "0+",
      "about": "The error code, or 0 if there was no error." },
    { "name": "Topics", "type": "[]TopicResult", "versions": "0",
      "about": "The write results for the topics.", "fields": [
      { "name": "Name", "type": "string", "versions": "0", "entityType": "topicName",
        "about": "The topic name." },
      { "name": "Partitions", "type": "[]PartitionResult", "versions": "0",
        "about": "The results for the partitions.", "fields": [
        { "name": "PartitionIndex", "type": "int32", "versions": "0",
          "about": "The partition index." },
        { "name": "ErrorCode", "type": "int16", "versions": "0",
          "about": "The error code, or 0 if there was no error." }
      ]}
    ]}
  ]
}

FindCoordinatorRequest

The FindCoordinatorRequest object is extended to support a new coordinator type:

public enum CoordinatorType { 
    GROUP((byte) 0), 
    TRANSACTION((byte) 1), 
    SHARE((byte) 2),
    MIRROR((byte) 3); // New type
}

Mirror Metadata Records

LastMirroredOffsets

LastMirroredOffsets record tracks the latest successfully mirrored offset for each partition.

{
  "apiKey": 1,
  "type": "coordinator-key",
  "name": "LastMirroredOffsetsKey",
  "validVersions": "0",
  "flexibleVersions": "none",
  "fields": [
    { "name": "MirrorName", "type": "string", "versions": "0",
      "about": "The cluster mirror name."}
  ]
}

{
  "apiKey": 1,
  "type": "coordinator-value",
  "name": "LastMirroredOffsetsValue",
  "validVersions": "0",
  "flexibleVersions": "0+",
  "fields": [
    { "name": "Topics", "type": "[]Topic", "versions": "0+",
      "about": "The mirror topics for which we want to store the last mirrored offsets.",  "fields": [
      { "name": "Name", "type": "string", "versions": "0",
        "about": "The topic name." },
      { "name": "Partitions", "type": "[]Partition", "versions": "0+",
        "about": "Each partition to record the last mirrored offsets.", "fields": [
        { "name": "PartitionIndex", "type": "int32", "versions": "0+",
          "about": "The partition index." },
        { "name": "LastMirroredOffset", "type": "int64", "versions": "0+",
          "about": "The last mirrored offset for this partition." }
      ]}
    ]}
  ]
}

MirrorPartitionState

MirrorPartitionState record represents the lifecycle states of a mirrored partition.

{
  "apiKey": 2,
  "type": "coordinator-key",
  "name": "MirrorPartitionStateKey",
  "validVersions": "0",
  "flexibleVersions": "none",
  "fields": [
    { "name": "MirrorName", "type": "string", "versions": "0",
      "about": "The cluster mirror name."}
  ]
}

{
  "apiKey": 2,
  "type": "coordinator-value",
  "name": "MirrorPartitionStateValue",
  "validVersions": "0",
  "flexibleVersions": "0+",
  "fields": [
    { "name": "TopicName", "type": "string", "versions": "0",
      "about": "The topic name."},
    { "name": "Partition", "type": "int32", "versions": "0",
      "about": "The partition index."},
    { "name": "State", "type": "int8", "versions": "0+",
      "about": "The mirror partition state." }
  ]
}

Configuration

A new configuration resource type is added for cluster mirrors, which is stored in the cluster metadata internal log:

public enum Type {
    // ... existing types ... 
    MIRROR((byte) 64, "mirror"); // New type
}


Cluster mirrors can be configured using the following properties:

Topic Configuration

Key

Description

Default

Dynamic

mirror.name

Identifies the mirror that manages this topic. Topics with this configuration set are read-only and can only be modified through mirror management APIs.

“”

yes

mirror.replication.throttled.replicas

A list of replicas for which log replication should be throttled on the mirror follower node. The list should describe a set of replicas in the form [PartitionId]:[BrokerId],[PartitionId]:[BrokerId]:... or alternatively the wildcard '*' can be used to throttle all replicas for this topic."

MAX_LONG

yes

Broker Configuration

Key

Description

Default

Dynamic

mirror.topic.num.partitions

Number of partitions for __mirror_state internal topic.

50

no

mirror.topic.replication.factor

Replication factor for __mirror_state internal topic. 

3

no

mirror.num.replica.fetchers

Number of fetcher threads per mirrored source broker,

1

yes

mirror.metadata.refresh.interval.ms

The interval in milliseconds at which the coordinator refreshes metadata from source clusters. This controls how frequently the coordinator polls source clusters to detect new topics and metadata changes.

30000

yes

request.timeout.ms

Request timeout for source cluster communication.

30000


socket.connection.setup.timeout.ms

Socket connection setup timeout.

10000


reconnect.backoff.ms

Backoff time before reconnection attempts. 

50


send.buffer.bytes

TCP send buffer size.

131072


receive.buffer.bytes

TCP receive buffer size.

65536


replica.fetch.backoff.ms

Time to wait before retrying fetch requests after failures (e.g., source leader change).



replica.fetch.max.bytes

Maximum bytes to fetch per partition in a single request to the source cluster.



replica.fetch.min.bytes

Minimum bytes that must be available before the source cluster responds to fetch requests (helps reduce cross-datacenter request frequency for low-throughput topics). 



replica.fetch.response.max.bytes

Maximum total bytes across all partitions in a single fetch response from source cluster (important for WAN bandwidth management in cluster mirroring).



replica.fetch.wait.max.ms

Maximum time the source cluster will wait to accumulate replica.fetch.min.bytes before responding (balances latency vs. efficiency for cross-cluster replication).



replica.socket.receive.buffer.bytes

TCP receive buffer size for connections to source cluster brokers (larger values can improve throughput over high-latency WAN links).



replica.socket.timeout.ms

Socket timeout for read operations from source cluster (should account for cross-datacenter network latency).



mirror.replication.throttled.rate

A long representing the upper bound (bytes/sec) on replication traffic for mirrored follower node enumerated in the property “mirror.replication.throttled.replicas” (for each topic). This property can be only set dynamically. It is suggested that the limit be kept above 1MB/s for accurate behaviour.


yes

Mirror Configuration

Key

Description

Default

Dynamic

bootstrap.servers

List of host/port pairs of the source cluster.



mirror.topic.properties.exclude

A comma-separated list of topic config property names to exclude from synchronization. Properties in this list will not be replicated from the source cluster. The mirror.name property is always excluded regardless of this setting.

follower.replication.throttled.replicas, leader.replication.throttled.replicas, message.timestamp.difference.max.ms, log.message.timestamp.before.max.ms, log.message.timestamp.after.max.ms, message.timestamp.type, unclean.leader.election.enable, min.insync.replicas, mirror.name

yes

mirror.groups.include

A comma-separated list of regex patterns for consumer group IDs to include in offset synchronization. Only consumer groups whose IDs match at least one of the patterns will have their offsets replicated from the source cluster.

.*

yes

mirror.acl.include

A comma-separated list of ACL include rules. Each rule uses semicolon-separated fields: resourceType;resourceName;operation;permissionType;principal. Use '*' as wildcard for any field. The resourceName field supports regex patterns. Trailing wildcard fields can be omitted. See AclRule javadoc for examples.

*

yes

security.protocol

Protocol for source cluster communication (PLAINTEXT, SSL, SASL_PLAINTEXT, SASL_SSL).



sasl.mechanism

SASL mechanism (PLAIN, SCRAM-SHA-256, SCRAM-SHA-512, GSSAPI, OAUTHBEARER).



sasl.jaas.config

JAAS login context parameters for authentication.



sasl.client.callback.handler.class

Fully qualified name of SASL client callback handler class.



sasl.login.callback.handler.class

Fully qualified name of SASL login callback handler class.



sasl.login.class

Fully qualified name of class implementing Login interface.



sasl.kerberos.service.name

Kerberos principal name for source cluster (when using GSSAPI).



sasl.kerberos.ticket.renew.jitter

Percentage of random jitter added to Kerberos ticket renewal time.



sasl.kerberos.ticket.renew.window.factor

Login thread sleep time until renewal as percentage of ticket lifetime.



sasl.kerberos.min.time.before.relogin 

Minimum time before attempting Kerberos credential renewal.



sasl.login.refresh.window.factor

Login refresh thread sleep factor relative to credential lifetime.



sasl.login.refresh.window.jitter

Maximum random jitter relative to credential refresh time.



sasl.login.refresh.min.period.seconds

Minimum time between credential refreshes.



sasl.login.refresh.buffer.seconds

Buffer time before credential expiration to maintain.



sasl.oauthbearer.token.endpoint.url

OAuth token endpoint URL (when using OAUTHBEARER).



sasl.oauthbearer.scope.claim.name

OAuth scope claim name for token requests.



sasl.oauthbearer.sub.claim.name

OAuth subject claim name for principal identification.



ssl.protocol

SSL protocol version (TLSv1.2, TLSv1.3).



ssl.provider

Name of security provider for SSL connections.



ssl.cipher.suites

List of enabled SSL cipher suites.



ssl.enabled.protocols

List of enabled SSL/TLS protocol versions.



ssl.keystore.type

Keystore file format (JKS, PKCS12, PEM).



ssl.keystore.location

Path to keystore file containing client certificate and private key.



ssl.keystore.password

Password for the keystore file.



ssl.keystore.key

Private key in PEM format (alternative to keystore file).



ssl.keystore.certificate.chain

Certificate chain in PEM format (alternative to keystore file).



ssl.key.password

Password for the private key in the keystore.



ssl.truststore.type

Truststore file format (JKS, PKCS12, PEM).



ssl.truststore.location

Path to truststore file for verifying source cluster broker certificates.



ssl.truststore.password

Path to truststore file for verifying source cluster broker certificates.



ssl.truststore.certificates

Trusted certificates in PEM format (alternative to truststore file).



ssl.keymanager.algorithm

Algorithm used by KeyManager factory (default: SunX509).



ssl.trustmanager.algorithm

Algorithm used by TrustManager factory (default: PKIX).



ssl.endpoint.identification.algorithm

Endpoint identification algorithm for hostname verification (https or empty to disable).



ssl.secure.random.implementation

SecureRandom PRNG implementation for SSL cryptography.



ssl.engine.factory.class

Fully qualified name of class implementing SslEngineFactory for custom SSL engine creation.



Metrics

A core set of metrics will be provided with the initial implementation.

Name

Type

Group

Tags

Description

JMX Bean

MaxLag

MirrorFetcherManager

kafka.server.mirror

clientId=MirrorReplica

Max lag in messages between destination leader and source leader replicas.

kafka.server.mirror:type=MirrorFetcherManager,name=MaxLag,clientId=MirrorReplica

MinFetchRate

MirrorFetcherManager

kafka.server.mirror

clientId=MirrorReplica

The min fetch rate between destination leader and source leader replicas.

kafka.server.mirror:type=MirrorFetcherManager,name=MirrorReplica

ConsumerLag

FetcherLagMetrics

kafka.server

clientId=MirrorFetcherThread-{sourceBroker.id}-{fetcherId}-{mirrorName},topic=([-.\w]+),partition=([0-9]+)

Lag in messages per remote leader replica.

kafka.serverr:type=FetcherLagMetrics,name=ConsumerLag,clientId=MirrorFetcherThread-{sourceBroker.id}-{fetcherId}-{mirrorName},topic=([-.\w]+),partition=([0-9]+)

DeadThreadCount

MirrorFetcherManager

kafka.server.mirror

clientId=MirrorReplica

Number of dead mirror fetcher threads.

kafka.server,mirror:type=MirrorFetcherManager,name=DeadThreadCount,clientId=MirrorReplica

FailedPartitionsCount

MirrorFetcherManager

kafka.server.mirror

clientId=MirrorReplica

Total count for failed partitions for any reason like auth, authorization, failed network with source.

kafka.serve.mirror:type=MirrorFetcherManager,name=FailedPartitionsCount,clientId=MirrorReplica

BytesPerSec

FetcherStats

kafka.server

clientId=MirrorFetcherThread-{sourceBroker.id}-{fetcherId}-{mirrorName},brokerHost={host},brokerPort={port}

Extend kafka.server.FetcherStats to report mirror fetcher threads.

kafka.server:type=FetcherStats,name=BytesPerSec,clientId=MirrorFetcherThread-{sourceBroker.id}-{fetcherId}-{mirrorName},brokerHost={host},brokerPort={port},mirror-name={mirrorName}

RequestsPerSec

FetcherStats

kafka.server

MirrorFetcherThread-{sourceBroker.id}-{fetcherId}-{mirrorName},brokerHost={host},brokerPort={port}

Extend kafka.server.FetcherStats to report mirror fetcher threads.

kafka.server:type=FetcherStats,name=RequestsPerSec,cclientId=MirrorFetcherThread-{sourceBroker.id}-{fetcherId}-{mirrorName}, brokerHost={host},brokerPort={port},mirror-name={mirrorName}



[LocalTimeMs,MessageConversionsTimeMs,

RemoteTimeMs,RequestBytes,

RequestQueueTimeMs,ResponseQueueTimeMs,

ResponseSendTimeMs,TemporaryMemoryBytes,

TotalTimeMs]

RequestMetrics

kafka.network

request=[mirror_requests]

Extend kafka.network:type=RequestMetrics to list cluster mirror requests.

kafka.network:type=RequestMetrics,name=*, request=*

ErrorsPerSec

RequestMetrics

kafka.network

request=[mirror_requests],error=*

Extend kafka.network:type=RequestMetrics to list cluster mirror requests.

kafka.network:type=RequestMetrics,name=ErrorsPerSec, request=*, error=*

RequestsPerSec

RequestMetrics

kafka.network

request=[mirror_requests],version=*

Extend kafka.network:type=RequestMetrics to list cluster mirror requests.

kafka.network:type=RequestMetrics,name=RequestsPerSec, request=*, version=*

connection-close-rate,

connection-close-total,

connection-count, connection-

creation-rate, connection-

creation-total, failed-authentication-rate, failed-authentication-total, failed-

reauthentication-rate, failed-

reauthentication-total,

incoming-byte-rate, incoming-byte-total, network-io-rate,

network-io-total, outgoing-

byte-rate, outgoing-byte-total,

reauthentication-latency-avg,

reauthentication-latency-max,

request-rate, request-size-avg,

request-size-max, request-total,

response-rate, response-total,

select-rate, select-total,

successful-authentication-no-

reauth-total, successful-

authentication-rate, successful-

authentication-total,

successful-reauthentication-

rate, successful-

reauthentication-total

mirror-broker-{DestinationBroker.id}-fetcher-{fetcherId}-mirror-{mirrorName}-metrics

kafka.server

broker-id={sourceBroker.id},fetcher-id={fetcherId}

Fetcher requests in the cluster mirror metrics.

kafka.server:type=mirror-broker-{sourceBroker.id}-fetcher-{fetcherId}-mirror-{mirrorName}-metrics,broker-id={sourceBroker.id},fetcher-id={fetcherId}

MetadataRefreshError

MirrorMetadataManager

kafka.server.mirror


Number of topic metadata refresh sync errors.

kafka.server.mirror:type=MirrorMetadataManager,name=aclSyncError

TopicConfigMetadataSyncError

MirrorMetadataManager

kafka.server.mirror


Number of topic configuration sync errors.


ConsumerGroupOffsetSyncError

MirrorMetadataManager

kafka.server.mirror


Number of CGs sync errors.


AclSyncError

MirrorMetadataManager

kafka.server.mirror


Number of ACLs sync errors.

kafka.server.mirror:type=MirrorMetadataManager,name=aclSyncError

byte-rate

MirrorReplication

kafka.server


Bandwidth quota metrics. Indicates the throttled data mirror replication rate of the broker in bytes/sec.

kafka.server:type=MirrorReplication

FailedPartitionState

MirrorMetadataManager

kafka.server.mirror


Number of partitions in failed state.

kafka.server.mirror:type=MirrorMetadataManager,name=FailedPartitionState

StoppedPartitionState

MirrorMetadataManager

kafka.server.mirror


Number of partitions in a stopped state.

kafka.server.mirror:type=MirrorMetadataManager,name=StoppedPartitionState

StoppingPartitionState

MirrorMetadataManager

kafka.server.mirror


Number of partitions in stopping state.

kafka.server.mirror:type=MirrorMetadataManager,name=StoppingPartitionState

MirroringPartitionState

MirrorMetadataManager

kafka.server.mirror


Number of partitions in mirroring state.

kafka.server.mirror:type=MirrorMetadataManager,name=MirroringPartitionState

PreparingPartitionState

MirrorMetadataManager

kafka.server.mirror


Number of partitions in preparing state.

kafka.server.mirror:type=MirrorMetadataManager,name=PreparingPartitionState

Compatibility, Deprecation, and Migration Plan

Cluster Mirroring will be introduced through a phased rollout across multiple Kafka releases to ensure stability and gather community feedback. 

Phases

Early access

Cluster Mirroring is introduced as an early access feature, disabled by default to prevent accidental production usage. To enable it, all cluster nodes (controllers and brokers) must explicitly enable unstable API versions and unstable feature versions in all configuration files. After starting the cluster with a minimum metadata version, operators can dynamically enable the mirror version feature to activate Cluster Mirroring. This stage is intended for testing and evaluation in non-production environments only, as the new APIs and metadata record formats may change in subsequent releases without backward compatibility guarantees.

Preview

In a future release, Cluster Mirroring will transition to preview status with frozen protocol and metadata schemas. The feature will still require explicit enablement via dynamic feature upgrades but will no longer require the unstable API and feature configuration. The feature remains disabled by default to ensure operators consciously opt-in, but the upgrade path from early access clusters will be officially supported with compatibility guarantees. This stage is suitable for pre-production testing and pilot deployments where API stability is required but production-grade maturity is not yet needed.

General availability

When Cluster Mirroring reaches general availability, the feature will be enabled by default when clusters reach the corresponding production metadata version. All new APIs will become stable production APIs with all unstable markers removed from their definition. No special configuration flags or explicit feature enablement will be required beyond setting an appropriate metadata version, and the feature will be fully supported for mission-critical production workloads under Kafka's standard compatibility guarantees. Clusters using Cluster Mirroring in preview can upgrade seamlessly to GA releases without migration steps. Downgrade is also supported, but it would require manual cleanup of the internal topic.

Migration From MirrorMaker 2

Cluster Mirror is not compatible with MirrorMaker 2. This is a critical consideration for users planning to migrate from MirrorMaker 2 to Cluster Mirroring.

MM2 and Cluster Mirroring use different internal topic structures and naming conventions for storing metadata and offsets. The two systems track and store consumer offsets differently, making it impossible to seamlessly transition between them.

Follow this process to switch from MirrorMaker 2 to Cluster Mirroring:

  1. Stop MM2 replication
  2. Delete mirror topics on destination cluster, including MM2 internal topics
  3. Start fresh with Cluster Mirroring

Compatibility Matrix

Note that some features require support from the source cluster.

Feature

Source Cluster Requirement

Destination Cluster Requirement

Notes

Core mirroring and failover

2.1

4.x

Kafka 4 is compatible with old clients versions up to 2.1 included.

Failback (reverse mirroring)

4.x

4.x

Requires last mirrored offset tracking on both sides, otherwise it will fallback and truncate to zero, effectively mirroring from scratch.

Share Groups

4.x

4.y

If the source doesn't support share groups, mirroring continues but share group offsets won't be synchronized.

Performance

MirrorFetcherThread uses the same fetch protocol optimizations as ReplicaFetcherThread:

Cluster Mirroring introduces additional replication threads and network I/O on brokers configured as read-only leaders for mirror partitions. The performance impact on existing intra-cluster replication is minimized through resource isolation:

Future Work

Synchronous mirroring: Currently, mirroring is asynchronous. The source cluster acknowledges the producer without waiting for the destination to replicate the data. Sync mirroring would guarantee that records are replicated to the destination cluster before the source acknowledges the produce request, providing stronger durability guarantees at the cost of higher latency. This would be useful for workloads where zero data loss across clusters is a strict requirement.

Tiered storage: Mirror topics in the destination cluster currently only replicate data from local storage on the source broker. Integrating with tiered storage would allow mirroring to handle data that has been offloaded to remote storage (e.g., S3, HDFS), enabling full replication of topics with long retention periods without requiring all data to reside in local broker storage.

Diskless topics: Diskless topics store data exclusively in tiered storage, with no local log segments on brokers. Supporting mirroring for diskless topics requires adapting the fetch and replication mechanisms to work without local storage, which introduces changes to how mirror offsets are tracked and how truncation is handled during failover.

Source cluster mirroring replication quota: Add source-side throttling that allows source cluster leaders to limit bandwidth served to all mirror fetchers, similar to how leader.replication.throttled.rate controls intra-cluster replication. This provides independent control over mirror catch-up traffic without impacting local replication or consumer workloads. Combined with destination-side throttling, operators gain complete bidirectional bandwidth control for mirror traffic.

Test Plan

Unit Tests

Unit tests will cover individual component behavior:

Integration Tests

Integration tests will validate end-to-end functionality across multiple brokers:

System Tests

System tests will validate behavior under realistic production conditions:

Rejected Alternatives

Keep Using MirrorMaker 2

This KIP introduces native cluster mirroring to address the limitations of MirrorMaker 2 described in the motivation section. The following tables provide a detailed comparison across deployment, features, and performance characteristics.

Deployment Comparison


MirrorMaker 2Cluster Mirroring
Architecture

External Connect workers (separate JVM)

Integrated into Kafka brokers using native replication protocol

Operational Complexity

Manage Connect cluster lifecycle independently

Unified with broker operations

Monitoring

Separate Connect metrics and dashboards

Standard Kafka JMX metrics

Feature Support Comparison


MirrorMaker 2Cluster Mirroring

Offset Translation

Lossy, requires remapping, causes reprocessing overhead

None needed - offsets preserved exactly

Metadata Sync

Requires separate connector configuration (MirrorSourceConnector, MirrorCheckpointConnector)

Automatic (topics, configs, consumer groups, ACLs)

Transactional Topics

Markers copied as regular records, incomplete transactions possible during replication

Markers mirrored, LSO truncation ensures consistency before failover

Topic Write Protection

Not supported (mirror topics always writable)

Read-only enforcement during mirroring, writable only after explicit failover

Tiered Storage

Fetches from broker (which reads from remote storage)

Not initially supported (future work)

Active-Active

Supported via topic prefixing and cycle detection

Not supported (read-only enforcement prevents cycles)

Share Groups

Not supported

Supported

Failback

Full re-mirror from offset 0

Delta sync using LastMirroredOffset

Topic Name Preservation

No - destination topics prefixed with source cluster alias (e.g., source.topic-name)

Yes - same topic name as source

Topic ID Preservation

No - destination gets new topic ID

Yes - same topic ID as source

Performance & Isolation Comparison


MirrorMaker 2Cluster Mirroring

Compression Overhead

Decompress + recompress records

Preserve source compression (zero overhead)

Failover Time

Offset translation + consumer group sync + bootstrap reconfiguration

Consumer group sync + bootstrap reconfiguration only

Bandwidth Control (Source)

None (unless client quotas manually configured)

Client quotas (current), dedicated mirror throttling (future)

Bandwidth Control (Destination)

None (unless client quotas manually)

Replica-level throttling (mirror.replication.throttled.rate)

Resource Isolation (Source)

Shares network bandwidth and disk I/O with all consumers

Shares network bandwidth and disk I/O with all consumers (same - current); isolated on replica-level (future)

Resource Isolation (Destination

Separate JVM heap (memory isolated from brokers) but still consume resources from the broker as any client

Shares network, disk I/O; dedicated thread pool within broker JVM

Malformed Batch Handling

Crashes Connect task, all partitions in task restart

Fails individual partition (FAILED state), others continue

Blast Radius(Failure)

All partitions in task affected

Single partition affected

Catch-Up Surge Protection

Connect worker heap may OOM, affects all tasks

Throttling on destination replica + source quota prevent memory spikes

Unclean Leader Election

No protection against log divergence

Explicitly unsupported with warning logs

Use Case Guidance:

Support Unclean Leader Election

Since there is no shared leader epoch between source and destination clusters, the destination cannot detect log divergence caused by unclean leader election on the source. With clean leader elections, the LastMirroredOffsets API is sufficient to reconcile the two clusters during failback. With unclean leader elections, it is not.

Clean leader election: (LMO works)

Consider a mirror where the destination has replicated offsets 0 and 1 from the source:

source (foo-0):               destination (foo-0):
offset 0, value: A            offset 0, value: A
offset 1, value: B            offset 1, value: B

A clean leadership change on the source bumps its epoch and appends a new record at offset 2. Before the destination fetches this record, the source goes down and failover occurs. The destination records LMO = 1 and becomes writable. A producer appends a different record at offset 2:

source (foo-0):               destination (foo-0):
offset 0, value: A            offset 0, value: A
offset 1, value: B            offset 1, value: B
offset 2, value: C            offset 2, value: D

When the old source comes back and wants to reverse-mirror from the destination, it queries the LastMirroredOffsets API and gets LMO = 1. It truncates its log to offset 1, then starts fetching from the destination. The logs converge.

Unclean leader election (LMO is not sufficient)

Now consider the same setup, but the leadership change on the source is an unclean election. The new source leader has lost data and starts with a divergent log:

source after unclean election (foo-0):           destination (foo-0):
offset 0, value: X                               offset 0, value: A
                                                 offset 1, value: B

The destination has not yet detected this divergence. The source appends more data. Before the destination can fetch and discover the epoch change, the source goes down and failover occurs. The destination records LMO = 1:

source (foo-0):               destination (foo-0):
offset 0, value: X            offset 0, value: A
offset 1, value: Y            offset 1, value: B
                              offset 2, value: D  (written after failover)

When the old source reverse-mirrors, it queries LMO = 1 and truncates to offset 1. But this only removes offset 1 onward, while offset 0 still contains X on the source vs A on the destination. The divergence at offsets below the LMO cannot be detected or resolved, because the destination has no record of the source's epoch history to compare against.

In summary, LMO-based truncation assumes the source log up to the LMO is a prefix of the destination log. Clean leader elections preserve this invariant. Unclean leader elections break it by allowing the source log to diverge at arbitrary offsets, including offsets already replicated to the destination. Resolving this would require a shared epoch mechanism or cross-cluster log reconciliation protocol, which is out of scope for this KIP.