DUE TO SPAM, SIGN-UP IS DISABLED. Goto Selfserve wiki signup and request an account.
...
Cluster mirrors can be configured using the following properties:
Topic Configuration
Key | Description | Default | Dynamic |
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 |
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 for source cluster communication. | 30000 | ||
Socket connection setup timeout. | 10000 | ||
Backoff time before reconnection attempts. | 50 | ||
send.buffer.bytes | TCP send buffer size. | 131072 | |
receive.buffer.bytes | TCP receive buffer size. | 65536 | |
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). | ||
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). | ||
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. | ||
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). | ||
OAuth scope claim name for token requests. | |||
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.
Metric 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.mirrorr: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
- What impact (if any) will there be on existing users?
- If we are changing behavior how will we phase out the older behavior?
- If we need special migration tools, describe them here.
- When will we remove the existing behavior?
Test Plan
Describe in few sentences how the KIP will be tested. We are mostly interested in system tests (since unit-tests are specific to implementation details). How will we know that the implementation works as expected? How will we know nothing broke?
Rejected Alternatives
...
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, administrators 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 administrators 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 (MM2). 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:
- Stop MM2 replication
- Delete mirrored topics on destination cluster, including MM2 internal topics
- 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. |
Tiered Storage | 3.0 | 4.y | If the source doesn't support Tiered Storage, mirroring continues but tiered segments won't be synchronized. |
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:
- Fetch Sessions: Incremental fetch sessions (KIP-227) reduce fetch request size by sending only changed partition metadata. This optimization is critical for cross-cluster replication where WAN latency is higher than LAN latency.
- Pipelining: Multiple fetch requests can be in-flight simultaneously, improving throughput over high-latency connections. The number of in-flight requests is controlled by standard replica fetcher settings.
- Compression: Record batches are transferred in their original compressed format, minimizing network bandwidth. The destination cluster decompresses and recompresses based on its own compression settings only if the compression codec differs.
- Zero-Copy Transfer: Within the destination cluster, replication from read-only leaders to followers uses zero-copy transfers where supported by the operating system.
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:
- Separate Thread Pools: Cross-cluster fetcher threads run in a dedicated thread pool, which is independent from the intra-cluster fetcher thread pool. This separation ensures that cross-cluster replication latency does not impact local replica synchronization.
- Network I/O Overhead: Read-only leaders perform additional network I/O to fetch from source clusters. This overhead is proportional to the number of mirror partitions and the replication throughput. Brokers with many mirror partitions may experience increased CPU usage for network processing and data serialization.
- Memory Footprint: Each mirror fetcher thread maintains its own fetch session state, partition state map, and response buffers. With default configuration, memory overhead is comparable to standard replica fetchers. The metadata manager maintains connection pools and metadata caches, adding minimal memory overhead.
- Bandwidth Consumption: Cross-cluster traffic between source and destination clusters consumes WAN bandwidth. For large-scale deployments, administrators should provision adequate inter-datacenter connectivity or configure throttling.
- State Management: Mirror partition state management is evenly distributed to available brokers to avoid any hot spot, especially during rolling update or restart events.
Test Plan
Unit Tests
Unit tests will cover individual component behavior:
- MirrorCoordinator: State loading, partition assignment, metadata persistence.
- MirrorMetadataManager: Topic creation, config sync, offset commit, ACL sync.
- MirrorFetcherThread: Epoch tracking, fetch processing, leader change handling.
- MirrorCommand: Command-line parsing, Admin API invocation, error handling.
- Protocol Serialization: Extended and new APIs serialization and deserialization.
Integration Tests
Integration tests will validate end-to-end functionality across multiple brokers:
- CLI Workflow: Create mirror with kafka-mirrors.sh, add topics, verify replication
- Basic Replication: Create mirror via API, replicate topic, verify data consistency
- Metadata Sync: Modify topic config in source, verify automatic sync to destination
- Partition Expansion: Add partitions to source topic, verify destination expands
- Consumer Groups: Commit offsets in source, verify replication to destination
- ACL Replication: Create ACL in source, verify creation in destination
- Leader Changes: Trigger leader election in source, verify fetcher reconnects
- Broker Failures: Stop destination broker, verify replication continues after recovery
System Tests
System tests will validate behavior under realistic production conditions:
- Performance Benchmark: Measure replication throughput and latency across WAN.
- Scalability Test: Replicate 1000 topics with 100,000 partitions across clusters.
- Failover Test: Simulate source cluster failure, measure consumer recovery time.
- Long-Running Stability: Run continuous replication for 7 days, verify no memory leaks or performance degradation.
- Security Validation: Test all authentication mechanisms (SASL PLAIN, SCRAM, Kerberos, mTLS) via kafka-mirrors.sh config files.
Rejected Alternatives
- Keep using MirrorMaker 2:
This KIP is to address the drawbacks existing MirrorMaker 2 as described in the motivation section.