04 / 05

A primary region fails. How would you fail over producers and consumers while controlling duplicates and gaps?

Difficulty: 8/10
Multi-region recovery, RPO/RTO, MirrorMaker

Failing Over Producers and Consumers with Controlled Duplicates and Gaps

When a primary region fails, the failover must handle three things: producers, consumers, and the data in flight. The first step is to declare the failure and decide to fail over. This decision should be based on health checks and should be made by a human or an automated system with a clear runbook. The second step is to redirect producers to the standby cluster. This is usually done by updating a DNS record or a configuration that the producers use to find the cluster. Producers must be configured to handle the switch: they need to know the new bootstrap servers and, if they use transactions or idempotence, they need to reset their state. Producers that were writing to the primary and did not receive an ack may have written records that were not replicated; those records will be lost unless the producer retries to the standby. The third step is to redirect consumers to the standby cluster. Consumers must be reconfigured with the new bootstrap servers, and their offsets must be translated from the primary to the standby using the offset mappings that MirrorMaker 2 maintains. Without offset translation, consumers would either start from the beginning (causing duplicates) or from the latest (causing gaps). The trade-off is between duplicates and gaps. If you resume from the last committed offset in the primary, you may lose the records that were in flight but not committed (a gap). If you resume from an earlier offset, you may reprocess records that were already processed (duplicates). The right choice depends on the business: for financial transactions, duplicates are usually worse than gaps, so you resume from the last committed offset and accept that some in-flight records may be lost. For other use cases, gaps may be worse, and you resume from an earlier offset to ensure completeness.

The mechanism for controlling duplicates and gaps is offset translation and idempotency. MirrorMaker 2 maintains an offset-syncs topic that maps source offsets to target offsets. When a consumer fails over, it can use this mapping to translate its committed offset from the primary to the standby. If the consumer was at offset 1000 in the primary and the mapping says that offset 1000 in the primary corresponds to offset 950 in the standby, the consumer resumes from 950. Records between 950 and 1000 in the standby may be duplicates of records the consumer already processed, or they may be records that were replicated after the consumer's last commit. To handle duplicates, the consumer must be idempotent. For producers, the failover must handle the case where a producer wrote to the primary but the record was not replicated. If the producer retries to the standby, it may create a duplicate if the record was actually replicated. Idempotent producers with transactional IDs can help, but the transactional state is per-cluster, so a producer that fails over needs a new transactional ID or a way to reconcile. The trade-off is between complexity and correctness. A simple failover that resumes from the last committed offset and relies on idempotency is easier to operate; a more complex failover that reconciles in-flight records gives better correctness but requires more tooling. Version note: MM2 offset translation is not perfect; it works for records that have been replicated, but it cannot account for records that were in flight at the time of failure. Always test the failover procedure with realistic failure scenarios and measure the actual duplicates and gaps.

A common mistake is to fail over without translating offsets, which causes either massive duplication or data loss. Another mistake is to assume that the standby cluster is fully caught up; replication lag means the standby is behind, and the gap is the RPO. A third mistake is to fail over producers and consumers at the same time without coordination; if producers start writing to the standby before consumers are ready, the consumers may miss records. The correct order is: declare failure, stop producers (or let them fail), wait for replication to catch up as much as possible, redirect consumers with translated offsets, then redirect producers. The trade-off is between speed and correctness. A fast failover reduces RTO but may increase duplicates or gaps; a slower, more careful failover reduces duplicates and gaps but increases RTO. For a financial platform, correctness is usually more important, so a slower, more careful failover is acceptable. Version note: some managed Kafka services provide automated failover with offset translation and idempotency guarantees. If you use a managed service, check its failover capabilities and test them. If you build your own, you need to build the offset translation, the consumer redirection, and the producer redirection, and you need to test them regularly.

javascript
  1. 1

    Failover must handle producers, consumers, and in-flight data.

  2. 2

    Redirect consumers using offset translation from MM2's offset-syncs topic.

  3. 3

    Resume from the last committed offset to avoid gaps, and rely on idempotency to handle duplicates.

  4. 4

    Redirect producers after consumers are ready to avoid missed records.

  5. 5

    Replication lag determines the gap; the standby is behind by the RPO.

  6. 6

    Use idempotent consumers and producers to handle duplicates during failover.

  7. 7

    Test the failover procedure regularly and measure actual duplicates and gaps.

  8. 8

    Managed services may provide automated failover; check and test their capabilities.

Share

Share via WhatsApp, X, Facebook, LinkedIn or copy link. Open Graph preview enabled.