skip to content

Walk through how a consumer application actually fails over to the standby cluster using translated checkpoints. What must the client do, and what are the pitfalls?

level: seniorimportance: should knowfreq 45%

answer

  1. Repoint bootstrap.servers → standby
  2. RemoteClusterUtils.translateOffsets / sync.group.offsets.enabled
  3. seek() before first poll, else auto.offset.reset bites
  4. IdentityReplicationPolicy keeps topic names
  5. Never double-active the same group while syncing
  6. Expect duplicates → idempotent consumers

basics

~20 s

The consumer reconnects to the standby's bootstrap servers, looks up its group's translated offsets (via RemoteClusterUtils/MirrorClient or pre-synced __consumer_offsets), seeks each partition to those offsets, then resumes. Pitfalls: stale checkpoints causing reprocessing, topic-name prefixes, and running the same group active on both clusters.

solid answer

~40 s

On failover the consumer must (1) change its bootstrap.servers to the standby, (2) obtain the target-cluster offsets for its group. If MM2 ran with sync.group.offsets.enabled=true, the standby's __consumer_offsets already holds translated commits, so the group can just subscribe and resume. Otherwise the app calls RemoteClusterUtils.translateOffsets or MirrorClient.remoteConsumerOffsets to read the checkpoints topic, then explicitly seek() each partition before polling. (3) It must use the correct topic names — under DefaultReplicationPolicy the standby topic is prefixed (primary.orders), so DR designs usually use IdentityReplicationPolicy to keep names identical. Pitfalls: translated offsets are approximate, so expect some reprocessing (design idempotent or at-least-once consumers); never let the same group consume actively on both clusters while offsets are being synced, or commits collide; and ensure the standby group isn't sync-overwritten after the live group starts committing there.

code

java · 13 lines
java
Map<String, Object> props = new HashMap<>();
props.put("bootstrap.servers", "standby-broker:9092");

// Translate the group's primary offsets to the standby cluster
Map<TopicPartition, OffsetAndMetadata> translated =
    RemoteClusterUtils.translateOffsets(props, "primary", "order-consumers", Duration.ofSeconds(30));

KafkaConsumer<String, String> consumer = new KafkaConsumer<>(consumerProps);
consumer.assign(translated.keySet());
for (Map.Entry<TopicPartition, OffsetAndMetadata> e : translated.entrySet()) {
    consumer.seek(e.getKey(), e.getValue().offset()); // seek BEFORE first poll
}
consumer.poll(Duration.ofMillis(500));

go deeper

for a junior

Know the consumer must point at the new cluster and resume from a translated position, not the same number.

for a middle

Explain the two ways to get translated offsets and the need to seek before polling.

for a senior

Cover IdentityReplicationPolicy, double-active hazards, and at-least-once reprocessing.

for a principal

Decide MM2 vs Cluster Linking, mandate idempotent consumers, and design config/DNS for transparent cutover.

## The consumer's job at failover A Kafka consumer resumes work based on **committed offsets** stored per consumer-group, per topic-partition. When the primary dies, those committed offsets live on the dead cluster. To continue on the standby, the consumer needs the **equivalent positions on the standby**, then must point itself there. ## Step-by-step 1. **Repoint connectivity.** Update `bootstrap.servers` (and security config) to the standby cluster. This usually comes from config/DNS rather than code, so DR-aware apps externalize the bootstrap endpoint. 2. **Obtain translated offsets.** Two mechanisms: - **Pre-synced**: if MM2 ran `MirrorCheckpointConnector` with `sync.group.offsets.enabled=true`, the standby's internal `__consumer_offsets` already contains the translated commits for inactive groups. The group can simply `subscribe()` and Kafka starts it at those positions. - **On-demand**: call `RemoteClusterUtils.translateOffsets(props, sourceClusterAlias, groupId, timeout)` or `MirrorClient.remoteConsumerOffsets(...)`. This returns a `Map<TopicPartition, OffsetAndMetadata>` of standby offsets derived from the checkpoints topic. The app then calls `consumer.seek(tp, offset)` for each partition **before** its first `poll()` that would otherwise honor `auto.offset.reset`. 3. **Use the right topic names.** Under the **DefaultReplicationPolicy**, the mirrored topic is renamed `<sourceAlias>.<topic>` (e.g. `primary.orders`). A consumer that subscribed to `orders` on the primary must now subscribe to `primary.orders` — awkward. DR setups therefore commonly use **IdentityReplicationPolicy** so the topic keeps its original name on the standby and the application config doesn't change. 4. **Start polling.** The consumer now reads from the translated positions. ## Pitfalls - **Reprocessing (at-least-once)**: translation rounds back to the nearest offset-sync, so the resume point is slightly behind the true last commit. Consumers must tolerate duplicates — design **idempotent** processing or dedup keys. Never assume exactly-once across the cutover. - **Stale checkpoints on unplanned failure**: if the primary crashed, the last few checkpoints may not have been emitted, widening the reprocessing window. - **Double-active group**: if the same `group.id` is consuming on **both** clusters while offset sync is on, the sync can overwrite the live group's commits (or vice versa), corrupting positions. Rule: only **inactive** groups should be offset-synced; once you fail a group over and it starts committing on the standby, stop syncing that group onto it. - **auto.offset.reset surprise**: if you seek incorrectly or skip the translation, a group with no offsets on the standby falls back to `auto.offset.reset` (earliest/latest), causing massive reprocessing or data skips. - **Transactions/EOS**: exactly-once consumers/streams have additional state (transactional IDs, changelog topics) that complicates a clean cross-cluster resume. ## Why teams sometimes pick Cluster Linking instead Cluster Linking mirrors topics **offset-preserving**, so the standby offsets equal the source offsets and consumer-offset migration is 1:1 — no `translateOffsets`, no approximate rounding. That removes much of the per-consumer failover complexity at the cost of being a Confluent feature.

  • What goes wrong if the standby has no committed offsets for the group and the app forgets to seek?
    The consumer falls back to auto.offset.reset. With 'earliest' it reprocesses the entire retained log; with 'latest' it skips everything produced before it joined, losing data. Either way it ignores the translated checkpoint, so you must seek() to the translated offsets first.
  • Why must you avoid running the same group.id actively on both clusters during offset sync?
    MM2's group-offset sync writes translated source offsets into the target's __consumer_offsets. If the group is also live and committing on the target, the sync overwrites those live commits (or races with them), corrupting positions. Only inactive groups should be synced.

saying these in an interview costs you the question

  • Assuming the consumer can keep the same numeric offsets on the standby
  • Forgetting to seek() and relying on auto.offset.reset
  • Running the same group active on both clusters while offsets are synced
  • Treating cross-cluster resume as exactly-once with no duplicates
  • Ignoring topic-name prefixing under DefaultReplicationPolicy

context