DUE TO SPAM, SIGN-UP IS DISABLED. Goto Selfserve wiki signup and request an account.
...
- State Management: Mirror configuration and partition states are stored in the internal topic. The coordinator loads the state on startup and partition leadership changes.
- Partition Assignment: Mirror partitions are distributed evenly to coordinators across all brokers and allows for horizontal scaling. The number of coordinator partitions is configurable via mirror.topic.num.partitions.
- Leader Election: When a broker becomes the leader for a __mirror_state partition, it loads the mirror metadata for all mirrors assigned to that partition and begins coordinating those mirrors. On resignation, it clears its in-memory state to avoid stale metadata.
- Metadata Refresh Scheduling: The coordinator schedules periodic metadata refresh operations by invoking a metadata manager every 30 seconds by default. This ensures that topology changes, configuration updates, and offset commits in the source cluster are continuously propagated to the destination cluster. The refresh interval is configurable via mirror.metadata.refresh.interval.ms.
- State Transitions: The coordinator manages asynchronous state transitions for mirror partitions. Each partition is an independent replication unit with its own state. When the coordinator is the leader for a mirror partition, it writes the state updates directly to the internal topic. Remote brokers read and write partition state via new RPCs, enabling distributed coordination across the cluster.
Figure 3: Mirror Partition Lifecycle.
...
| Code Block | ||
|---|---|---|
| ||
public enum Type {
// ... existing types ...
MIRROR((byte) 64, "mirror"); // New type
} |
Cluster mirrors can be configured using the following properties:
...
Each configuration or group of configurations in the following list has a specific scope (mirror, topic, broker).
Mirror Configuration
Set via CreateMirror or IncrementalAlterConfigs. Stored in cluster metadata records.
Key | Description | Default |
|---|
bootstrap.servers | A list of host/port pairs to use for establishing the initial connection to the source cluster. | |
mirror. |
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
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 |
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
, 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 | ||
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. | .* |
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. | * |
security.protocol | Protocol for source cluster communication (PLAINTEXT, SSL, SASL_PLAINTEXT, SASL_SSL). | |
sasl.* | SASL configuration properties. | |
ssl.* | SSL configuration properties. |
Topic Configuration
Set by topic creation or alter. Stored in topic config.
Key | Description | Default |
|---|---|---|
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. | |
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 |
Broker Configuration
Set via broker config. Stored in server.properties or dynamic broker config.
Key | Description | Default |
|---|---|---|
mirror.topic.num.partitions | Number of partitions for __mirror_state internal topic. | 50 |
mirror.topic.replication.factor | Replication factor for __mirror_state internal topic. | 3 |
mirror.num.replica.fetchers | Number of fetcher threads per mirrored source broker, | 1 |
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 |
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. | |
request.timeout.ms | Maximum amount of time in milliseconds the client will wait for the response of a request. | 30000 |
socket.* | Socket connection configurations. | |
replica.* | Fetcher threads configurations |
...
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
. |
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
...
