10.4 Disaster recovery planning
That is a genuinely expensive operation to discover mid-incident.
4.1 RTO and RPO
- RTO — Recovery Time Objective. “the Maximum amount of time before all services must resume after a disaster”
- ► “The Lower the RTO, the more important it is to Avoid manual processes and application restarts, since Very low RTO can be achieved only with automated failover.”
- RPO — Recovery Point Objective. “the Maximum amount of time for which data may be lost as a result of a disaster”
- ► “Low RPO requires Real-time mirroring with low latencies, and RPO = 0 requires synchronous replication.”
The implication chain: RPO=0 ⇒ synchronous replication ⇒ you cannot use MirrorMaker ⇒ you need a stretch cluster (§6) or a commercial synchronous solution (§9).
4.2 ⚠️ Data loss and inconsistencies in unplanned failover
The arithmetic — do this for your own throughput:
“Because Kafka’s various mirroring solutions are ALL ASYNCHRONOUS, the DR cluster WILL NOT HAVE THE LATEST MESSAGES from the primary cluster.”
Worked example:
1,000,000 messages/second
× 5 ms lag between primary and DR
═════════════════════════════════
= 5,000 messages behind IN THE BEST-CASE SCENARIO“in a busy system you should expect the DR cluster to be A FEW HUNDRED OR EVEN A FEW THOUSAND MESSAGES BEHIND the primary.” ► “PREPARE FOR UNPLANNED FAILOVER TO INCLUDE SOME DATA LOSS.”
Planned failover is different: "you can stop the primary cluster and WAIT for the mirroring process to mirror the remaining messages before failing over applications to the DR cluster, thus avoiding this data loss."
And the cross-topic consistency problem:
"note that mirroring solutions currently DON'T SUPPORT TRANSACTIONS, which means that if some events in multiple topics are related to each other (e.g., SALES and LINE ITEMS), you can have SOME events arrive to the DR site in time for the failover and OTHERS THAT DON'T. Your applications will need to be able to handle A LINE ITEM WITHOUT A CORRESPONDING SALE after you failover."
(This is exactly Ch. 8 §3.6's finding: MirrorMaker can be per-record exactly-once but cannot preserve transaction atomicity.)
4.3 The four failover-offset strategies
"One of the challenging tasks in failing over to another cluster is making sure applications know where to start consuming data."
Strategy A: Auto offset reset — simplest, lossiest
No offset mirroring at all. Pick one auto.offset.reset:
| Setting | What happens on failover |
|---|---|
earliest | “start reading from the beginning of available data and HANDLE LARGE AMOUNTS OF DUPLICATES” |
latest | “skip to the end and MISS AN UNKNOWN (and hopefully small) NUMBER OF EVENTS” |
“If your application handles duplicates with no issues, or missing some data is no big deal, THIS OPTION IS BY FAR THE EASIEST. Simply SKIPPING TO THE END of the topic on failover is A POPULAR FAILOVER METHOD DUE TO ITS SIMPLICITY.”
Strategy B: Replicate the __consumer_offsets topic — three serious caveats
"If you mirror this topic to your DR cluster, when consumers start consuming from the DR cluster, they will be able to pick up their old offsets and continue from where they left off. It is simple, but THERE IS A LONG LIST OF CAVEATS."
- ⚠ CAVEAT 2 — PRODUCER RETRIES CAUSE DIVERGENCE. “even if you started mirroring immediately when the topic was created and BOTH start with 0, PRODUCER RETRIES CAN CAUSE OFFSETS TO DIVERGE.”
- ⚠ CAVEAT 3 — OFFSET COMMITS AND RECORDS ARRIVE OUT OF STEP. “because of the LAG between primary and DR clusters AND because mirroring solutions DON’T SUPPORT TRANSACTIONS, an offset committed by a Kafka consumer MAY ARRIVE AHEAD OR BEHIND THE RECORD WITH THIS OFFSET. A consumer that fails over MAY FIND COMMITTED OFFSETS WITHOUT MATCHING RECORDS. Or it may find that the latest committed offset in the DR site is OLDER than the latest committed offset in the primary.”
What you must decide in advance:
- Accept DUPLICATES if the DR’s latest committed offset is older than the primary’s, or if DR record offsets are ahead due to retries.
- Decide what to do when “the latest committed offset in the DR site DOESN’T HAVE A MATCHING RECORD — do you start processing from the BEGINNING of the topic or SKIP TO THE END?”
"this approach has its limitations. Still, this option lets you failover with a REDUCED NUMBER of duplicated or missing events compared to other approaches while still being simple to implement."
Strategy C: Time-based failover — 💡 the good compromise
| Available since | What it added |
|---|---|
| 0.10.0 | “each message includes a Timestamp indicating the time the message was sent to Kafka” |
| 0.10.1.0 | “brokers include An index and an API for looking up offsets By the timestamp” |
“if you failover to the DR cluster and you know that Your trouble started at 4:05 a.m., you can tell consumers to start processing data From 4:03 a.m. There will be some duplicates from those two minutes, but it is probably better than other alternatives”
And the human argument, which is unusually candid and quite persuasive:
"the behavior is MUCH EASIER TO EXPLAIN TO EVERYONE IN THE COMPANY — 'We failed back to 4:03 a.m.' sounds better than 'We failed back to what may or may not be the latest committed offsets.'"
Two ways to implement it:
- Bake it into your app. “Have a user-configurable option to specify the Start time for the app. If configured, the app can use the new APIs to Fetch offset by time, Seek to that time, and start consuming from the right point, committing offsets as usual.”
- ► “great If you wrote All your applications this way In advance. But what if you didn’t?”
kafka-consumer-groupsTOOL (timestamp-based reset added in 0.11.0)- ⚠ “The consumer group Should be stopped while running this type of tool and Started immediately after.”
bin/kafka-consumer-groups.sh --bootstrap-server localhost:9092 \
--reset-offsets --all-topics --group my-group \
--to-datetime 2021-03-31T04:03:00.000 --execute"This option is recommended in deployments that need to guarantee a level of certainty in their failover."
(Note the "stop the group first" requirement — same constraint as Ch. 5 §6.4's alterConsumerGroupOffsets.)
Strategy D: Offset translation — most precise
“Offsets are stored WHENEVER THE DIFFERENCE BETWEEN THE TWO OFFSETS CHANGES.”
- ► “There is NO NEED TO STORE ALL THE OFFSET MAPPINGS between 495 and 596; we just ASSUME THAT THE DIFFERENCE REMAINS THE SAME.”
- THE HISTORY: “some organizations chose to use an EXTERNAL DATA STORE, such as Apache Cassandra, to store mapping of offsets from one cluster to another.” TODAY: “mirroring solutions, INCLUDING MIRRORMAKER, USE A KAFKA TOPIC for storing offset translation metadata.”
- THEN at failover: “instead of mapping TIMESTAMPS (which are always a bit inaccurate) to offsets, WE MAP PRIMARY OFFSETS TO DR OFFSETS and use those.” Use strategy A or C’s mechanism to force consumers onto them.
- ⚠ “This STILL has an issue with offset commits that ARRIVED AHEAD of the records themselves and offset commits that DIDN’T GET MIRRORED to the DR on time, but IT COVERS SOME CASES.”
Summary of the four:
| Strategy | Precision | Complexity | Residual problem |
|---|---|---|---|
| A. auto.offset.reset | Worst | Trivial | Mass duplicates or unknown data loss |
B. Replicate __consumer_offsets | Poor–medium | Low | Offsets diverge; commits ≠ records |
| C. Time-based | Good | Low–medium | Bounded duplicates; explainable |
| D. Offset translation | Best | Medium (built into MM2) | Commits still race the records |
4.4 ⚠️ After the failover — you probably have to scrape the old primary
"It is tempting to simply modify the mirroring processes to reverse their direction and start mirroring from the new primary to the old one. However, this leads to two important questions:"
- “How do we know where to start mirroring? We need to solve The same problem we have for all our consumers — for the mirroring application itself. And remember that all our solutions have cases where they either cause duplicates or miss data — Sometimes both.”
- “it is likely that Your original primary will have events that the dr cluster does not. If you just start mirroring new data back, The extra history will remain and the two clusters will be inconsistent.”
F and G never made it to the DR cluster; H, I and J are new writes since failover. Reverse-mirroring naively puts F and G back on the old primary alone.
The remedy: "for scenarios where consistency and ordering guarantees are critical, the simplest solution is to FIRST SCRAPE THE ORIGINAL CLUSTER — DELETE ALL THE DATA AND COMMITTED OFFSETS — and then start mirroring from the new primary back to what is now the new DR cluster. This gives you A CLEAN SLATE that is identical to the new primary."
That is a genuinely expensive operation to discover mid-incident. Plan and rehearse it.
4.5 Cluster discovery
"in the event of failover, your applications will need to know how to start communicating with the failover cluster. If you HARDCODED THE HOSTNAMES of your primary cluster brokers in the producer and consumer properties, THIS WILL BE CHALLENGING."
The common approach: “create a DNS name that usually points to the primary brokers. In case of an emergency, the DNS name can be Pointed to the standby cluster.”
- 💡 “The discovery service (DNS or other) Doesn’t need to include all the brokers — Kafka clients only need to access A single broker successfully in order to get metadata about the cluster and discover the others. So, Including just three brokers is usually fine.”
- ⚠ “Regardless of the discovery method, Most failover scenarios do require bouncing consumer applications after failover So they can find the new offsets from which they need to start consuming.”
- ► “For Automated failover Without application restart to achieve Very low RTO, Failover logic should be built into client applications.”
(Note the client.dns.lookup interaction from Ch. 5 §3.1 — a DNS alias plus SASL needs resolve_canonical_bootstrap_servers_only.)