Two Clusters, One Change Feed: Replicating Across a WAN Without Stretching Zaris
Sooner or later someone asks the question: we have a cluster in one region and users in another — can we just add a few nodes over there and let one cluster span both? It is a reasonable instinct. One cluster, one namespace, one connection string, data everywhere. The instinct is also a trap, and it is worth being precise about why.
A Zaris cluster is engineered for a datacenter network. Membership is tracked with heartbeats on short timers. Ownership of a partition moves between nodes through a handshake that assumes the other side answers quickly. Failover fences a suspected-dead owner and promotes a replica on the order of seconds. Every one of those mechanisms is calibrated for a link where the round trip is measured in microseconds and a "slow" node is genuinely sick — not merely eighty milliseconds and an ocean away. Put half the nodes across a WAN and you haven't built a geo-distributed cluster; you've built one cluster that now mistakes normal WAN latency for failure, flaps ownership, and risks split-brain the first time the link between regions hiccups.
So the right shape is not one stretched cluster. It is two independent clusters — each self-contained, each with its own ownership and failover staying entirely on its own LAN — joined by an asynchronous, resumable change feed. This post is about that pattern: why independence is the load-bearing decision, and how Zaris's operation log already hands you the one thing a WAN replicator actually needs — a cursor you can resume from.
Why a stretched cluster is the wrong tool
It helps to name the specific assumptions a single cluster makes, because each one breaks differently over a WAN.
- Heartbeat-based membership decides a node is gone when it misses a few beats. A cross-region link that adds latency and the occasional partition will trip that detector during healthy operation. The cluster will declare live nodes dead.
- The ownership protocol moves partitions between nodes with a handshake that expects a prompt reply. Across a WAN, those transitions stall, retry, and compete — exactly the churn that the ownership and cutover machinery is designed to resolve quickly on a LAN.
- Failover timers promote a replica when the owner looks dead. Over a WAN, "looks dead" and "is just far away" become indistinguishable, so you get spurious promotions — and if the link heals with a promoted replica on each side, you have two owners for one partition.
None of this is a defect. It is a store doing exactly what a low-latency in-datacenter store should do, applied to a network it was never meant to span. The fix is not to loosen every timer until the cluster tolerates a WAN — that would also make it tolerate real failures far too slowly. The fix is to stop asking one cluster to do it.
The building block you already have: a replayable log
Zaris keeps, per segment, a revision-ordered operation log. Every accepted write advances a monotonic revision and appends an entry. That log is the substrate replicas and new partition owners use to catch up without freezing writes: a consumer that knows the last revision it saw can replay forward from that point and converge to current truth.
That property — replay forward from a known revision — is the entire foundation of a WAN replicator. A change feed is nothing more than a long-lived reader of that log that streams entries in revision order and remembers where it stopped. The thing it remembers is a cursor: the highest revision it has durably applied, per segment.
Because the log is ordered by revision and replay starts from a point you choose, the feed has the two traits a WAN link demands:
- Resumable. The link will drop — that is the defining fact of a WAN. When it comes back, the follower reconnects and asks for everything after its stored cursor. No full re-copy, no coordination, just "resume from revision N."
- Idempotent on replay. Apply an entry only when its revision is greater than the last one applied for that segment. If a reconnect re-delivers the tail, the follower skips what it already has. Duplicate delivery becomes a no-op, which is what lets the feed be at-least-once instead of having to be exactly-once.
The pattern, end to end
Two clusters. Call them A (the source) and B (the follower). Each is an ordinary, independent Zaris cluster that knows nothing about the other's internals.
In cluster B you run a replication follower — a client that consumes A's change feed and writes the operations into B locally. Its state is tiny and durable: a cursor per segment. The loop is unglamorous, which is the point:
- Connect to A, request the change feed for the segments (or key prefixes) you want to replicate, starting after each segment's stored cursor.
- Receive entries in revision order. For each, apply the operation to B —
SET,DELETE, expiry — guarded by the revision check so re-delivery is harmless. - Advance and persist the cursor as entries are durably applied.
- On any disconnect, go back to step 1. The cursor is where you left off.
Everything the follower needs to survive interruption lives in that cursor. There is no distributed transaction across the WAN, no shared membership, no cross-region quorum. A hiccup on the link costs you nothing but a little freshness.
What you are trading: RPO is not zero, and that is the deal
This replication is asynchronous. The follower is always some distance behind the source — however many revisions have been produced in A but not yet applied in B. That distance is your Recovery Point Objective: if region A vanishes this instant, you lose whatever had not yet reached B.
That is not a flaw to apologize for; it is the trade you are explicitly choosing. The alternative — making writes in A wait for an acknowledgement from B before they commit — would put a WAN round trip on your write path and surrender the low latency that is the whole reason you run Zaris. Synchronous cross-region replication turns every write into a transcontinental request. Almost nobody actually wants that once they see the latency bill.
So the honest framing is: you pick the RPO. A well-fed follower on a healthy link stays within a handful of revisions — sub-second staleness in practice. A saturated or flapping link lets it drift further. The cursor makes the drift bounded and visible — you can measure exactly how far behind B is by comparing its cursor to A's current revision — but it does not make it zero. Size your tolerance accordingly, and alert on cursor lag the way you would alert on replica lag inside a cluster.
When the source fails over, the feed does not care
Here is where building on the operation log pays off a second time. Inside cluster A, the partition your feed is reading will occasionally change owners — a node dies, a replica is promoted, a partition migrates. The log survives that: it is keyed by revision, and a new owner continues the same revision sequence. For the follower, a source-side failover is just another reconnect. It re-establishes the feed against A's current topology and resumes from its cursor. The revisions it already applied are still the revisions it already applied; the ones it is missing are still after the cursor. The follower never has to understand who in A owns the partition — only what revision it last saw.
This is the quiet reason the cursor-over-log design is worth more than a naive "tail the writes and ship them" relay. A relay that streams live writes with no durable, revision-keyed position has nothing to resume from when the source topology shifts. The log gives the follower a stable coordinate system that outlives any individual owner.
Scope is a choice, and conflicts are a boundary
Two more things to be deliberate about.
Data scope. You do not have to replicate everything. The feed is per-segment, and you can select the segments or key prefixes that matter — a specific tenant, a specific dataset, the slice that a disaster-recovery region actually needs. Replicating less means less WAN traffic and a smaller follower; replicating everything gives you a full warm standby. Decide on purpose rather than by default.
Direction and conflicts. The pattern above is one-directional: A is the source of truth, B is a read-mostly follower / standby. That is the clean case, and it has no conflict problem — there is exactly one writer per key. The moment you make it bidirectional — both regions accepting writes for the same keys — you have introduced concurrent writers separated by a WAN, and no amount of change-feed plumbing resolves that for you. You now owe a conflict-resolution policy: last-writer-wins on a clock you trust, per-key ownership so the two regions never write the same key, or application-level merge. Active-active is a real and sometimes necessary design, but it is a data-modeling decision, not a transport feature. Reach for one-directional first; adopt bidirectional only when you have decided, explicitly, how two writers reconcile.
The shape to remember
A WAN is not a slow LAN; it is a different medium with partitions as a normal event rather than an emergency. A single cluster, correctly, refuses to pretend otherwise — its timers and its ownership protocol are tuned to treat a long silence as a failure, because inside a datacenter it almost always is. So you don't stretch the cluster. You run two of them, each sovereign on its own network, and you join them with the one primitive that is built to tolerate a link that comes and goes: a revision-ordered log you can replay, and a cursor that remembers where you were.
Independence is the decision that makes the rest safe. The change feed is just the pipe. The operation log already gave you the cursor.