When would you implement a custom ReplicationPolicy, and what methods must it correctly implement to remain loop-safe and failover-correct?
answer
- implements ReplicationPolicy
- formatRemoteTopic <-> topicSource/upstreamTopic
- isInternalTopic
- must be invertible (round-trip)
- drives cycle detection + offset translation
basics
~20 sYou write a custom ReplicationPolicy when neither default prefixing nor identity naming fits — e.g. a bespoke naming scheme or org convention. It must correctly implement formatRemoteTopic plus the inverse parsing (topicSource, upstreamTopic) and isInternalTopic, or cycle detection and offset translation break.
solid answer
~40 sImplement a custom `ReplicationPolicy` (set via `replication.policy.class`) when you need a naming scheme the two built-ins can't express: e.g. embedding a region into a suffix, mapping to an existing legacy convention, routing some topics with prefixes and others without, or namespacing for multi-tenant clusters. The interface's contract is bidirectional, so you must keep the WRITE and READ sides perfectly inverse: `formatRemoteTopic(sourceAlias, topic)` builds the remote name; `topicSource(remoteTopic)` recovers the immediate source alias; `upstreamTopic(remoteTopic)` strips one hop back toward the original; and `isInternalTopic(topic)` flags topics MM2 must not treat as user data. If formatRemoteTopic and topicSource/upstreamTopic disagree, MM2 loses origin recovery — cycle detection in active/active fails and offset/checkpoint translation for failover produces wrong mappings. You should also handle multi-hop nesting and ensure your naming can't collide with real topic names or internal topics.
go deeper
Know that the policy class is pluggable via replication.policy.class.
Name the two built-ins and that a custom one is rarely needed.
List the key methods and that they must be mutually consistent for cycle detection and offset sync.
Reason about invertibility, multi-hop, collisions, internal-topic classification, blast radius, and when a custom policy is justified versus the built-ins.
## The ReplicationPolicy interface MM2 abstracts all topic-naming behind `org.apache.kafka.connect.mirror.ReplicationPolicy`, selected with `replication.policy.class`. The built-ins are `DefaultReplicationPolicy` (source-prefixed) and `IdentityReplicationPolicy` (unchanged names). A custom policy plugs into the same slot. ## When a custom policy is justified - **Bespoke naming conventions**: your org mandates `orders.us-west` (suffix) or `tenantA/orders` style namespacing that neither built-in produces. - **Selective prefixing**: prefix some topics (multi-source) but leave a curated set unprefixed (shared reference data) for seamless client failover. - **Legacy interop**: match a naming scheme an existing pre-MM2 mirroring setup established, so consumers don't change. - **Multi-tenancy**: encode tenant/region routing into names while staying loop-safe. Note: if you just want unprefixed names for active/passive, use IdentityReplicationPolicy — don't write a custom policy. ## The methods and their contract The interface is fundamentally about an invertible transform between (sourceAlias, originalTopic) and a remote name: - `formatRemoteTopic(sourceClusterAlias, topic)` — WRITE side: produce the replicated topic name. - `topicSource(topic)` — READ side: given a remote name, return the immediate source alias (or null if it's not a remote/replicated name). - `upstreamTopic(topic)` — READ side: strip one hop, returning the name as it was on the upstream cluster (peeling nested prefixes for multi-hop). - `isInternalTopic(topic)` — classify MM2/Kafka internal topics (offsets, heartbeats, checkpoints, transaction state) so they're filtered, not replicated as data. There are also default helper methods (e.g. for offset-sync, checkpoints, heartbeats topic names) you can override. ## Why correctness is unforgiving MM2 relies on these being mutually consistent: 1. **Cycle detection**: `MirrorSourceConnector` uses topicSource/upstreamTopic to decide a topic's origin and skip replicating it back to a cluster in its path. A buggy parser that can't recover origin reintroduces infinite loops in active/active. 2. **Offset and checkpoint translation**: `MirrorCheckpointConnector` maps consumer offsets on remote topics back to the source topic so consumers fail over to the right position. If upstreamTopic is wrong, consumers resume at wrong offsets (data loss or reprocessing). 3. **Internal-topic filtering**: if isInternalTopic misclassifies, you might replicate `__consumer_offsets` or MM2's own heartbeats as user data, creating chaos. ## Design pitfalls to avoid - **Non-invertible naming**: any scheme where formatRemoteTopic isn't reversible by topicSource/upstreamTopic is broken. Test round-trips, including multi-hop (A→B→C). - **Separator/name ambiguity**: if real topic names can contain your delimiter, parsing misfires; choose delimiters carefully or escape. - **Collisions**: ensure two distinct (alias, topic) pairs can't map to the same remote name. - **Forgetting internal topics**: default-classify the standard internal names plus your own. ## Configuration ``` replication.policy.class=com.example.RegionSuffixReplicationPolicy replication.policy.separator=. # honored only if your policy reads it ``` Your policy can read `replication.policy.separator` (and any custom configs) via the `configure(Map)` method if it implements `Configurable`. ## Bottom line A custom policy is a small class with an outsized blast radius: it governs cycle safety, failover correctness, and internal-topic hygiene simultaneously. Reach for it only when the built-ins genuinely can't express your scheme, and unit-test the write/read round-trip exhaustively.
- What's the single most important invariant a custom ReplicationPolicy must satisfy?The naming transform must be invertible: formatRemoteTopic and topicSource/upstreamTopic must be exact inverses (including across multiple hops), or MM2 loses origin recovery and both cycle detection and offset translation break.
- If you only need unprefixed names for active/passive DR, should you write a custom policy?No — use the built-in IdentityReplicationPolicy. A custom policy is only warranted when neither built-in can express your required naming scheme.
- Which connector relies on upstreamTopic for failover correctness?MirrorCheckpointConnector (offset/checkpoint translation) maps remote-topic consumer offsets back to the source topic using upstreamTopic; a wrong implementation makes consumers resume at incorrect offsets.
saying these in an interview costs you the question
- Writing a custom policy when IdentityReplicationPolicy already suffices
- Implementing formatRemoteTopic but neglecting topicSource/upstreamTopic (breaks origin recovery)
- Forgetting isInternalTopic, leading to internal topics being replicated as data
- Assuming a non-invertible naming scheme is acceptable