Learn Labs
10. Cross-Cluster Data Mirroring

10.6 Apache Kafka's MirrorMaker

① "Define aliases for the clusters used in replication flows."

6.1 The evolution — why MM1 was replaced

Legacy mirrormaker (mm1). “a collection of consumers that were Members of a consumer group to read data from a set of source topics, and a Shared kafka producer in each MirrorMaker process to send those events to the destination cluster.”

  • ⚠ “it had several issues, particularly Latency spikes as Configuration changes and Addition of new topics resulted in Stop-the-world rebalances.”

Mirrormaker 2.0 (introduced in Kafka 2.4.0). “the next-generation multicluster mirroring solution… Based on the kafka connect framework, overcoming many of the shortcomings of its predecessor. Complex topologies can be easily configured to support a wide range of use cases like Disaster recovery, Backup, Migration, and Data aggregation.”

MORE ABOUT MIRRORMAKER

"MirrorMaker sounds very simple, but because we were trying to be very efficient and get very close to exactly-once delivery, IT TURNED OUT TO BE TRICKY TO IMPLEMENT CORRECTLY. MirrorMaker has been REWRITTEN MULTIPLE TIMES."

6.2 Architecture

consumeproduceassignssource clusterMirrorMaker taskone consumer + one producertarget clusterKafka Connectassigns tasks to workers · REST API
  • MIRRORMAKER 2.0 = a Kafka Connect SOURCE CONNECTOR whose “database” is ANOTHER KAFKA CLUSTER.
  • The Connect framework assigns tasks to worker nodes: “you may have MULTIPLE TASKS ON ONE SERVER or have the tasks SPREAD OUT to multiple servers”.
  • ► “This REPLACES THE MANUAL WORK of figuring out how many MirrorMaker STREAMS should run per instance and how many INSTANCES per machine.”
  • ► Connect’s REST API centrally manages connector/task config.
Figure 10.6.16.2 Architecture

A nice operational argument: "If we assume that most Kafka deployments include Kafka Connect for other reasons (sending database change events into Kafka is a very popular use case), then by running MirrorMaker INSIDE Connect, we can CUT DOWN ON THE NUMBER OF CLUSTERS WE NEED TO MANAGE."

💡 The key MM2 design decisions

  1. NO CONSUMER GROUP PROTOCOL. “MirrorMaker ALLOCATES PARTITIONS TO TASKS EVENLY WITHOUT USING KAFKA’S CONSUMER GROUP-MANAGEMENT PROTOCOL to AVOID LATENCY SPIKES DUE TO REBALANCES when new topics or partitions are added.” ► This is the direct fix for MM1’s central flaw.
  2. PARTITION-PRESERVING. “Events from each partition in the source cluster are mirrored to THE SAME PARTITION in the target cluster, PRESERVING SEMANTIC PARTITIONING and MAINTAINING ORDERING OF EVENTS FOR EACH PARTITION.” “If new partitions are added to source topics, THEY ARE AUTOMATICALLY CREATED IN THE TARGET TOPIC.”
  3. MORE THAN JUST DATA — “a COMPLETE mirroring solution”:
    • consumer OFFSETS migration
    • topic CONFIGURATION migration
    • topic ACLs migration

Vocabulary: "A replication flow defines the configuration of a directional flow from a source cluster to a target cluster. Multiple replication flows can be defined to define COMPLEX TOPOLOGIES, including hub-and-spoke, active-standby, and active-active."

6.3 Configuring MirrorMaker

bin/connect-mirror-maker.sh etc/kafka/connect-mirror-maker.properties
Active-standby: New York → London
clusters = NYC, LON                                  # ① aliases

NYC.bootstrap.servers = kafka.nyc.example.com:9092   # ② prefix = alias
LON.bootstrap.servers = kafka.lon.example.com:9092

NYC->LON.enabled = true                              # ③ enable the flow
NYC->LON.topics = .*                                 # ④ what to mirror

① "Define aliases for the clusters used in replication flows." ② "Configure bootstrap for each cluster, using the cluster alias as the prefix." ③ "Enable replication flow using the prefix source->target. All configuration options for this flow use the same prefix." ④ The topics to mirror.

Mirror topics — the naming strategy that prevents cycles

"for each replication flow, a REGULAR EXPRESSION may be specified for the topic names... In this example, we chose to replicate every topic, but it is often good practice to use something like prod.* and AVOID REPLICATING TEST TOPICS. A separate topic EXCLUSION list containing names or patterns like test.* may also be specified."

⚠ “TARGET TOPIC NAMES ARE AUTOMATICALLY PREFIXED WITH THE SOURCE CLUSTER ALIAS BY DEFAULT.”

mirrorNYCtopic ordersLONtopic NYC.orders
  • WHY: “This default naming strategy PREVENTS REPLICATION CYCLES resulting in events being ENDLESSLY MIRRORED between the two clusters in active-active mode if topics are mirrored from NYC to LON as well as LON to NYC.”
  • BONUS: “The distinction between LOCAL and REMOTE topics also SUPPORTS AGGREGATION USE CASES since consumers may choose SUBSCRIPTION PATTERNS to consume data produced from JUST THE LOCAL REGION or subscribe to topics from ALL REGIONS to get the complete dataset.”
Figure 10.6.3Mirror topics — the naming strategy that prevents cycles

Auto-discovery: "MirrorMaker periodically checks for new topics in the source cluster and starts mirroring automatically if they match the configured patterns. If more partitions are added to the source topic, the SAME NUMBER of partitions is automatically added to the target topic, ensuring that events in the source topic appear in the same partitions in the same order in the target topic."

Consumer offset migration
  • RemoteClusterUtils — a utility class “to enable consumers to SEEK TO THE LAST CHECKPOINTED OFFSET in a DR cluster WITH OFFSET TRANSLATION when failing over”
  • Kafka 2.7.0+ — “PERIODIC MIGRATION of consumer offsets… to AUTOMATICALLY COMMIT TRANSLATED OFFSETS to the target __consumer_offsets topic so that consumers switching to a DR cluster can RESTART FROM WHERE THEY LEFT OFF in the primary cluster with NO DATA LOSS AND MINIMAL DUPLICATE PROCESSING.”
  • Which groups get migrated is CUSTOMIZABLE.
  • 💡 SAFETY INTERLOCK: “for added protection, MirrorMaker DOES NOT OVERWRITE OFFSETS IF CONSUMERS ON THE TARGET CLUSTER ARE ACTIVELY USING THE TARGET CONSUMER GROUP, thus avoiding any accidental conflicts.”

That interlock is the same principle as Ch. 5 §6.4 — never write offsets under a live group.

Topic configuration and ACL migration

Config migration: enabled by default with “reasonable periodic refresh intervals that may be sufficient in most cases.”

  • ⚠ “Most of the topic configuration settings from the source are applied to the target topic, BUT A FEW LIKE min.insync.replicas ARE NOT APPLIED BY DEFAULT. The list of excluded configs can be customized.”

ACL migration:

  • ⚠ “Only literal topic ACLs that match topics being mirrored are migrated, so if you are using Prefixed or wildcard ACLs or Alternative authorization mechanisms, you will need to configure those on the target cluster explicitly.”
  • ⚠ “ACLs for Topic:Write Are not migrated to ensure that Only mirrormaker is allowed to write to the target topic. appropriate access must be explicitly granted at the time of failover to ensure that applications work with the secondary cluster.”

Two failover landmines hidden in these defaults:

  1. min.insync.replicas is not migrated → your DR topics may have weaker durability than you assume.
  2. Topic:Write ACLs are deliberately not migrated → on failover, your producers cannot write until someone grants access. Put this in the runbook.
tasks.max

"limits the maximum number of tasks... The default is 1, but A MINIMUM OF 2 IS RECOMMENDED. When replicating a lot of topic partitions, HIGHER VALUES SHOULD BE USED if possible to increase parallelism."

Configuration prefixes — the hierarchy

“Kafka Connect and connector configs can be specified Without any prefix.”

“more specific prefixed configuration has Higher precedence than the less specific or nonprefixed configuration”

  • {cluster}.{connector_config}
  • {cluster}.admin.{admin_config}
  • {source_cluster}.consumer.{consumer_config}
  • {target_cluster}.producer.{producer_config}
  • {source_cluster}->{target_cluster}.{replication_flow_config}

6.4 Multicluster topologies

Active-active NYC ↔ LON — just enable both directions:

clusters = NYC, LON
NYC.bootstrap.servers = kafka.nyc.example.com:9092
LON.bootstrap.servers = kafka.lon.example.com:9092
NYC->LON.enabled = true
NYC->LON.topics  = .*
LON->NYC.enabled = true
LON->NYC.topics  = .*

"even though all topics from NYC are mirrored to LON and vice versa, MirrorMaker ensures that the same event isn't constantly mirrored back and forth since REMOTE TOPICS USE THE CLUSTER ALIAS AS THE PREFIX."

💡 Operational best practice: "It is good practice to use THE SAME CONFIGURATION FILE that contains the FULL replication topology for DIFFERENT MirrorMaker processes since it avoids conflicts when configs are shared using the internal configs topic in the target datacenter." Start processes with the --clusters option to specify which target cluster this process serves.

Fan out — add a third cluster:

clusters = NYC, LON, SF
SF.bootstrap.servers = kafka.sf.example.com:9092
NYC->SF.enabled = true
NYC->SF.topics  = .*

6.5 Securing MirrorMaker

"For production clusters, it is important to ensure that ALL cross-datacenter traffic is secure... SSL should be used to ENCRYPT ALL cross-datacenter traffic."

NYC.security.protocol=SASL_SSL          # ① match the broker listener
NYC.sasl.mechanism=PLAIN
NYC.sasl.jaas.config=org.apache.kafka.common.security.plain.PlainLoginModule \
    required username="MirrorMaker" password="MirrorMaker-password";   # ②

① "Security protocol should match that of the broker listener corresponding to the bootstrap servers specified for the cluster. SSL or SASL_SSL is recommended." ② "For SSL, keystores should be specified if mutual client authentication is enabled."

The complete ACL list MirrorMaker needs
SOURCE clusterTARGET cluster
Topic:Read (consume)Topic:Create + Topic:Write (create and produce)
Topic:DescribeConfigs (get source topic configuration)Topic:AlterConfigs (update target topic config)
—Topic:Alter (ADD PARTITIONS if new source partitions are detected)
Group:Describe (source consumer group metadata, including offsets)Group:Read (commit offsets for those groups)
Cluster:Describe (source topic ACLs)Cluster:Alter (update target topic ACLs)

BOTH: Topic:Create + Topic:Write for INTERNAL MirrorMaker topics.

Note how the ACL list maps 1:1 to the feature list — each capability (data, configs, partitions, offsets, ACLs) needs its own grant pair. Miss one and that one feature silently stops working.


On this page