Versions Compared

Key

  • This line was added.
  • This line was removed.
  • Formatting was changed.
Comment: formatting

...

  • Potential Data Loss: In the event of a catastrophic failure of the source cluster, recently produced records that have not yet been replicated to the destination cluster will be lost. The amount of data loss depends on replication lag at the time of failure.
  • RPO (Recovery Point Objective): Organizations must handle a non-zero RPO determined by the replication lag between source and destination clusters. Typical replication lag ranges from seconds to minutes depending on network bandwidth, throughput, and geographic distance.

 

Asynchronous replication should provide the right balance for disaster recovery 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.

...

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.

...

CreateMirrorResponse

Code Block


  "apiKey": TBD, 

  "type": "response", 

  "name": "CreateMirrorResponse", 

  "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." } 

  ] 

}

...

AddTopicsToMirrorRequest

Code Block
{

  "apiKey":TBD,

  "type": "request",

  "listeners": ["broker", "controller"],

  "name": "AddTopicsToMirrorRequest",

  "validVersions": "0",

  "flexibleVersions": "0+",

  "fields": [

    { "name": "Topics", "type": "[]TopicState", "versions": "0+", "about": "The topic state.",

      "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+", "nullableVersions": "0+",

          "about": "The mirror name."}

      ]}

  ]

}

AddTopicsToMirrorResponse

Code Block
{

  "apiKey":TBD,

  "type": "response",

  "name": "AddTopicsToMirrorResponse",

  "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": "TopicId", "type": "uuid", "versions": "0+", "about": "The unique topic ID."},

    { "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." }

  ]

}

RemoveTopicsFromMirror

RemoveTopicsFromMirrorRequest

Code Block
{

  "apiKey": TBD,

  "type": "request",

  "listeners": ["broker", "controller"],

  "name": "RemoveTopicsFromMirrorRequest",

  "validVersions": "0",

  "flexibleVersions": "0+",

  "fields": [

    { "name": "MirrorName", "type": "string", "versions": "0+", "ignorable": true,

      "about": "The cluster mirror name." },

    { "name": "Topics", "type": "[]TopicState", "versions": "0+", "about": "The topic state.",

      "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

Code Block
{

  "apiKey":96 TBD,

  "type": "response",

  "name": "RemoveTopicsFromMirrorResponse",

  "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": "TopicId", "type": "uuid", "versions": "0+", "about": "The unique topic ID."},

    { "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." }

  ]

}

LastMirroredOffset

ListMirroredOffsetsRequest

Code Block
{

  "apiKey":97 TBD,

  "type": "request",

  "listeners": ["broker", "controller"],

  "name": "LastMirroredOffsetsRequest",

  "validVersions": "0",

  "flexibleVersions": "0+",

  "fields": [

    { "name": "MirrorName", "type": "string", "versions": "0+", "about": "The mirror name." },

    { "name": "Topics", "type": "[]TopicState", "versions": "0",

      "about": "The responses per topic.", "fields": [

      { "name": "Name", "type": "string", "versions": "0", "entityType": "topicName",

        "about": "The topic name." },

      { "name": "Partitions", "type": "[]PartitionState", "versions": "0",

        "about": "The responses per partition.", "fields": [

        { "name": "PartitionIndex", "type": "int32", "versions": "0",

          "about": "The partition index." }

      ]}

    ]}

  ]

}

LastMirroredOffsetsResponse

Code Block
{

  "apiKey":97 TBD,

  "type": "response",

  "name": "LastMirroredOffsetsResponse",

  "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": "[]OffsetResponseTopic", "versions": "0",

      "about": "The responses per topic.", "fields": [

      { "name": "Name", "type": "string", "versions": "0", "entityType": "topicName",

        "about": "The topic name." },

      { "name": "Partitions", "type": "[]OffsetResponsePartition", "versions": "0",

        "about": "The responses per partition.", "fields": [

        { "name": "PartitionIndex", "type": "int32", "versions": "0",

          "about": "The partition index." },

        { "name": "LastMirroredOffset", "type": "int64", "versions": "0",

          "about": "The last mirrored record offset." },

        { "name": "ErrorCode", "type": "int16", "versions": "0",

          "about": "The error code, or 0 if there was no error." }

      ]}

    ]}

  ]

}

ListMirrors

ListMirrorsRequest

Code Block
{

  "apiKey": 98TBD,

  "type": "request",

  "listeners": ["broker"],

  "name": "ListMirrorsRequest",

  // Version 0 is the initial version.

  "validVersions": "0",

  "flexibleVersions": "0+",

  "fields": []

}

ListMirrorsResponse

Code Block
{

  "apiKey": 98TBD,

  "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+",

        "about": "The 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

DescribeMirrorsRequest

Code Block
{

  "apiKey": 99TBD,

  "type": "request",

  "listeners": ["broker"],

  "name": "DescribeMirrorsRequest",

  // Version 0 is the initial version.

  "validVersions": "0",

  "flexibleVersions": "0+",

  "fields": [

    { "name": "MirrorNames", "type": "[]string", "versions": "0+",

      "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

Code Block
{

  "apiKey": 99TBD,

  "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+",

        "about": "The 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

ReadMirrorStatesRequest

Code Block
{

  "apiKey":100 TBD,

  "type": "request",

  "listeners": ["broker", "controller"],

  "name": "ReadMirrorStatesRequest",

  "validVersions": "0",

  "flexibleVersions": "0+",

  "fields": [

    { "name": "MirrorName", "type": "string", "versions": "0+", "about": "The mirror name." },

    { "name": "Topics", "type": "[]TopicState", "versions": "0",

      "about": "The responses per topic.", "fields": [

      { "name": "Name", "type": "string", "versions": "0", "entityType": "topicName",

        "about": "The topic name." },

      { "name": "Partitions", "type": "[]PartitionState", "versions": "0",

        "about": "The responses per partition.", "fields": [

        { "name": "PartitionIndex", "type": "int32", "versions": "0",

          "about": "The partition index." }

        ]}

      ]},

    { "name": "NeedPartitionStates", "type": "bool", "versions": "0+", "default": "true",

      "about": "Need the partition states or only topics states needed." }

  ]

}

ReadMirrorStatesResponse

Code Block
{

  "apiKey":100 TBD,

  "type": "response",

  "name": "ReadMirrorStatesResponse",

  "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": "[]TopicState", "versions": "0",

      "about": "The responses per topic.", "fields": [

      { "name": "Name", "type": "string", "versions": "0", "entityType": "topicName",

        "about": "The topic name." },

      { "name": "Partitions", "type": "[]PartitionState", "versions": "0",

        "about": "The responses per partition.", "fields": [

        { "name": "PartitionIndex", "type": "int32", "versions": "0",

          "about": "The partition index." },

        { "name": "LastMirroredOffset", "type": "int64", "versions": "0",

          "about": "The last mirrored record 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

WriteMirrorStatesRequest

Code Block
{

  "apiKey":101 TBD,

  "type": "request",

  "listeners": ["broker", "controller"],

  "name": "WriteMirrorStatesRequest",

  "validVersions": "0",

  "flexibleVersions": "0+",

  "fields": [

    { "name": "MirrorName", "type": "string", "versions": "0+", "about": "The mirror name." },

    { "name": "TopicsUpdated", "type": "[]TopicState", "versions": "0",

      "about": "The topics to be updated.", "fields": [

      { "name": "Name", "type": "string", "versions": "0", "entityType": "topicName",

        "about": "The topic name." },

      { "name": "Partitions", "type": "[]PartitionState", "versions": "0",

        "about": "The responses per partition.", "fields": [

        { "name": "PartitionIndex", "type": "int32", "versions": "0",

          "about": "The partition index." },

        { "name": "LastMirroredOffset", "type": "int64", "versions": "0",

          "about": "The last mirrored record 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

Code Block
{

  "apiKey":101 TBD,

  "type": "response",

  "name": "WriteMirrorStatesResponse",

  "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": "[]TopicState", "versions": "0",

      "about": "The responses per topic.", "fields": [

      { "name": "Name", "type": "string", "versions": "0", "entityType": "topicName",

        "about": "The topic name." },

      { "name": "Partitions", "type": "[]PartitionState", "versions": "0",

        "about": "The responses per partition.", "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:

Code Block
languagejava
public enum CoordinatorType 

...


    GROUP((byte) 0), 

...



    TRANSACTION((byte) 1), 

...



    MIRROR((byte) 2); // New type

...


}

Mirror Metadata Records

LastMirroredOffsets

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

Code Block
{

...


  "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.

Code Block
{

...


  "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:

Code Block
languagejava
public enum Type {

...


    // ... existing types ... 

...


    MIRROR((byte) 64, "mirror"); // New type

...


}


Cluster mirrors can be configured using the following properties:

...

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.



...

  • 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:

...