10.7 Deploying MirrorMaker in production
A canary is the only listed monitor that measures the thing you actually care about — end-to-end delivery — rather than a proxy for it.
7.1 Deployment modes
DEDICATED MODE (what the example uses): “You can start ANY NUMBER of these processes to form A DEDICATED MIRRORMAKER CLUSTER that is SCALABLE AND FAULT-TOLERANT. The processes mirroring to the same cluster WILL FIND EACH OTHER AND BALANCE LOAD BETWEEN THEM AUTOMATICALLY.”
Production hygiene:
- run as a SERVICE, “in the background with
nohupand redirecting its console output to a log file” - or use the
-daemoncommand-line option - “Most companies that use MirrorMaker have their own STARTUP SCRIPTS”
- Ansible, Puppet, Chef, Salt “are often used to automate deployment and manage the many configuration options”
- “MirrorMaker may also be run inside a DOCKER CONTAINER.”
💡 “MirrorMaker is COMPLETELY STATELESS and DOESN’T REQUIRE ANY DISK STORAGE (all the data and state are stored in Kafka itself).”
All Connect deployment modes work:
| Mode | Use |
|---|---|
| Standalone | "development and testing" — one worker, one machine |
| A connector in an existing distributed Connect cluster | "by explicitly configuring the connectors" |
| Distributed — dedicated MM cluster or shared Connect cluster | ✅ "For production use, we RECOMMEND running MirrorMaker in DISTRIBUTED MODE" |
7.2 ⚠️ Where to run MirrorMaker — the most important operational decision
DEFAULT: “If at all possible, RUN MIRRORMAKER AT THE TARGET DATACENTER.” So NYC → SF means MirrorMaker runs IN SF and consumes across the US.
The reasoning, in full — this is the chapter's best paragraph:
*"long-distance networks can be a bit less reliable than those inside a datacenter. If there is a network partition and you lose connectivity between the datacenters, having a CONSUMER that is unable to connect to a cluster is MUCH SAFER than a PRODUCER that can't connect.
If the CONSUMER can't connect, it simply won't be able to read events, but THE EVENTS WILL STILL BE STORED IN THE SOURCE KAFKA CLUSTER and can remain there for a long time. THERE IS NO RISK OF LOSING EVENTS.
On the other hand, if the events were ALREADY CONSUMED and MirrorMaker CAN'T PRODUCE them due to network partition, THERE IS ALWAYS A RISK THAT THESE EVENTS WILL ACCIDENTALLY GET LOST by MirrorMaker.
So, REMOTE CONSUMING IS SAFER THAN REMOTE PRODUCING."*
Two exceptions where you must produce remotely
Exception 1: encryption asymmetry — and the zero-copy reason
*"When do you have to consume locally and produce remotely? The answer is when you need to encrypt the data while it is transferred between the datacenters but you don't need to encrypt the data inside the datacenter.
Consumers take a SIGNIFICANT PERFORMANCE HIT when connecting to Kafka with SSL encryption — MUCH MORE SO THAN PRODUCERS. This is because use of SSL requires COPYING DATA FOR ENCRYPTION, which means CONSUMERS NO LONGER ENJOY THE PERFORMANCE BENEFITS OF THE USUAL ZERO-COPY OPTIMIZATION. And this performance hit ALSO AFFECTS THE KAFKA BROKERS THEMSELVES."*
Why the asymmetry (ties back to Ch. 6 §5.4):
- Producers were never zero-copy, so SSL costs them comparatively less.
- ► THEREFORE: place MirrorMaker at the SOURCE datacenter, “having it consume UNENCRYPTED data locally, and then producing it to the remote datacenter through an SSL ENCRYPTED CONNECTION.”
⚠️ If you do this, you must compensate for the loss of the fail-safe:
“make sure MirrorMaker’s Connect producer is configured to Never lose events by configuring it with:
acks=all- a sufficient number of retries
Also, configure MirrorMaker to Fail fast using errors.tolerance=none when it fails to send events, Which is typically safer to do than to continue and risk data loss.”
💡 "Note that newer versions of Java have SIGNIFICANTLY INCREASED SSL PERFORMANCE, so producing locally and consuming remotely may be a viable option EVEN WITH ENCRYPTION." — i.e. re-measure before accepting the trade.
Exception 2: firewalls (hybrid cloud)
"a hybrid scenario when mirroring from an on-premises cluster to a cloud cluster. Secure on-premises clusters are likely to be BEHIND A FIREWALL THAT DOESN'T ALLOW INCOMING CONNECTIONS FROM THE CLOUD. Running MirrorMaker ON PREMISE allows ALL CONNECTIONS TO BE FROM ON PREMISES TO THE CLOUD."
7.3 Monitoring MirrorMaker
Kafka Connect metrics
"connector metrics to monitor connector status, source connector metrics to monitor throughput, and worker metrics to monitor rebalance delays. Connect also provides a REST API to view and manage connectors."
MirrorMaker-specific metrics
| Metric | Meaning |
|---|---|
replication-latency-ms | "the time interval between the record timestamp and the time at which the record was successfully produced to the target cluster. ... useful to detect if the target is not keeping up." |
record-age-ms | "the age of records at the time of replication" |
byte-rate | replication throughput |
checkpoint-latency-ms | "offset migration latency" |
| heartbeats | "MirrorMaker also emits periodic heartbeats by default, which can be used to monitor its health." |
How to interpret latency: "Increased latency during peak hours may be OK if there is sufficient capacity to catch up later, but SUSTAINED INCREASE in latency may indicate INSUFFICIENT CAPACITY."
⚠️ Lag monitoring — two methods, neither perfect
THE LAG: “the difference in offsets between the LATEST MESSAGE IN THE SOURCE cluster and the LATEST MESSAGE IN THE TARGET cluster.”
Example: source last offset = 7, target last offset = 5 → lag = 2.
| Method | Reports here | Why it is not accurate |
|---|---|---|
1 — the offset MirrorMaker COMMITTED to the source. Tool: kafka-consumer-groups | 5 (real lag is 2) | ⚠ “NOT 100% ACCURATE because MIRRORMAKER DOESN’T COMMIT OFFSETS ALL THE TIME. It commits offsets EVERY MINUTE by default, so you will see THE LAG GROW FOR A MINUTE AND THEN SUDDENLY DROP.” |
| 2 — the latest offset MirrorMaker READ. “The consumers embedded in MirrorMaker publish key metrics in JMX. One of them is the CONSUMER MAXIMUM LAG (over all the partitions it is consuming).” | 1 — it already read message 6, even though message 6 wasn’t produced to the destination yet | ⚠ “ALSO NOT 100% ACCURATE because it is updated based on WHAT THE CONSUMER READ but DOESN’T TAKE INTO ACCOUNT whether the PRODUCER MANAGED TO SEND those messages to the destination and whether they were ACKNOWLEDGED successfully.” |
► LinkedIn’s BURROW “monitors the same information but has A MORE SOPHISTICATED METHOD to determine whether the lag represents a real problem, SO YOU WON’T GET FALSE ALERTS.”
⚠ THE MONITORING GAP: “if MirrorMaker SKIPS OR DROPS MESSAGES, NEITHER METHOD WILL DETECT AN ISSUE because THEY JUST TRACK THE LATEST OFFSET.” ► Confluent Control Center “monitors MESSAGE COUNTS AND CHECKSUMS and CLOSES THIS MONITORING GAP.”
Producer and consumer metrics for tuning
| Side | Metrics |
|---|---|
| Consumer | fetch-size-avg, fetch-size-max, fetch-rate, fetch-throttle-time-avg, fetch-throttle-time-max |
| Producer | batch-size-avg, batch-size-max, requests-in-flight, record-retry-rate |
| Both | io-ratio, io-wait-ratio |
Canary
"If you monitor everything else, a canary isn't strictly necessary, but we like to add it in for MULTIPLE LAYERS OF MONITORING. It provides a process that, every minute, sends an event to a special topic in the source cluster and tries to read the event from the destination cluster. It also alerts you if the event takes more than an acceptable amount of time to arrive. This can mean that MirrorMaker is lagging or that it isn't available at all."
A canary is the only listed monitor that measures the thing you actually care about — end-to-end delivery — rather than a proxy for it.