Set up disaster recovery¶
Keep a live copy of a topic on a recovery cluster, usually in another region, and know how much you would lose if the first cluster were gone.
New in v3.2.0.
Before you start: every item of Before you start on both clusters, in both directions; the CLI with a context for each (a for the cluster in use, b for the recovery cluster); and a metrics store outside a's region.
A remote child of the topic on a sends every record to the same topic on b, at least once. b then holds everything a's consumers have not processed, up to the link's lag. When a's region is lost, you fail over to b; when it is back, you fail back. The link is an HTTPS client of b's public API, so it assumes no shared network: everything below holds between regions as within one.
Set up in peacetime¶
-
Create the topic on
bwitha's partition count and schema, under the same name, so that clients switch by URL alone. Size its retention as in Retention. -
Check the bounds in both directions.
b's host and port must be insideremotes.allowed_hostsandremotes.allowed_portson every node ofa, anda's inside them on every node ofb, since widening them is a config change and a restart. Every node needs egress to the other cluster's ingress. Name each cluster by a hostname pinned to it and its region, never a global failover name: failing over is your decision, not DNS's. -
Provision both directions now, so nobody creates a credential during an incident:
- the remote
bona, with the userrepl-from-a-7f3k9qonb; - the remote
aonb, with the userrepl-from-b-2m8x4dona.
Each cluster keeps its remotes in its own metadata, encrypted under its own cluster secret, so failing back never depends on
a's region. The two clusters never share a cluster secret (operating condition 5). A remote with no link sends nothing; rotate the idle pair on the same schedule as the live one. - the remote
-
Attach the link:
narad --ctx a topic attach orders orders-dr --remote b \ --from unconsumed --lanes 2unconsumedis enough for a queue: after a failover,b's consumers need only whatahas not processed.earliestsendsa's whole retained log at once (at 5 MB/s and 72 hours, 1.3 TB of egress, then a catch-up that competes with live traffic); use it only whenbmust hold the history. Chooselanesas in Throughput. -
Deploy
b's consumers at zero replicas, ready to scale. -
Set up the alerts in Watch the link, evaluated outside
a's region. -
Rehearse once. Pause the link for an hour (
narad --ctx a topic pause orders orders-dr --reason drill), watch the headroom fall, resume it, and time the catch-up.
Retention¶
- The source.
a's log is the only buffer while the link is down. A link down forDwith lagLwhen it stopped loses nothing whilea's retention exceedsL + Dplus the time to notice and fix whatever stopped it. Set at least 72 hours on a topic with a cross-region remote child: regional incidents have lasted most of a day, and 72 hours covers one that runs into a weekend, plus the catch-up. The floor is 24 hours, and the attach warns below 72. During an outage, raising the retention works at once while headroom is above 0. - Catching up. With the link's capacity
Cand the produce rateP, the oldest unshipped record gets younger byC/P - 1seconds every second once the link is back, so catching up takesD × P / (C - P): as long as the outage atC = 2P, five times as long atC = 1.2P. WithC ≤ Pthe link falls behind with no outage at all; the headroom alert fires before records age out. - What drop-behind leaves. If the oldest unshipped records age out of
aanyway, the cursor skips them and counts them innarad_fanout_child_dropped_messagesona.bgets no marker, and nothing can fill the gap later. - The recovery cluster. After a failover,
b's consumers start atb's oldest retained record, because nothing was ever acked there. Sob's retention must exceed the age ofa's oldest unprocessed record at the moment of failure, plus the time to fail over.bcounts retention from each record's arrival onb, so it keeps a record up to the link's lag longer thanadoes; where data has a mandated maximum retention, setb's short enough to meet it at the largest lag you tolerate.
Throughput¶
A lane sends one request at a time and waits for the answer, so its rate is bounded by the round trip: a request carries up to 1,000 records to a target on this release (100 to an older one), and up to 960 KiB. A key always travels on one lane, so one hot key can never go faster than one lane. More lanes per parent partition (up to 8) add parallel streams; the remote's max_in_flight limit (16 per node by default) caps them all. With remotes.allowed_hosts set, narad remote test reports each node's connect time and an estimate of one lane's records per second at that round trip. Size the lanes so the link's capacity is at least twice the produce rate.
The link sends JSON payloads as they are and anything else as base64. Between regions, egress is billed per GB by the source region, plus any NAT gateway on the path; prefer a private interconnect (peering, a transit gateway, a VPN), which narrows who can connect but does not replace TLS or the credential. compression: zstd on the remote (Change a remote) compresses each request that shrinks by at least 10%, against a target on this release; leave it off for a topic whose records must not leak to each other through compressed lengths.
Watch the link¶
At the moment a's region fails, the data at risk is the link's lag (records committed on a that b has not answered 202 for) plus a's ingress dispatch backlog of orders (records a answered 202 for and had not committed yet, which the lag cannot see). No metric counts that backlog for one topic: narad_ingress_dispatch_backlog_records counts every topic the node serves, so it is an upper bound. Nothing on b's side is at risk: a 202 from b is durable on b.
| Signal | Why | Alert |
|---|---|---|
max(narad_ |
The recovery point | Page above your objective |
narad_ on a |
An upper bound on the part of the recovery point the lag cannot see; node-wide, across every topic | Warn when it keeps rising for 10 minutes: records are not reaching their owners. A node that takes produce is rarely at 0, so do not alert on above 0 |
sum(rate(narad_ |
Below 1, the link is falling behind | Warn below 1 for 15 minutes |
narad_ |
Time left before drop-behind | Warn below 12 hours, page below 4 |
narad_ |
Why the link stopped | As in Monitor and alert |
narad_ |
Below 960 KiB, the path is cutting uploads short | Warn at the 64 KiB floor for 10 minutes |
- Watch
afrom outsidea's region. Senda's metrics, over an authenticated connection, to a store inb's region or a global one. Do not opena's metrics listener to another region: it serves without credentials, for scrapes inside the cluster. After a failure, the data at risk is then the last storedlag_secondsplus one scrape interval, plus at most the last dispatch backlog. - Optionally, watch from
b's side too. Give a small topic onaits own remote child tob, produce{"ts": <a's Unix ms>}to it every 5 seconds, and page frombwhen the newest one is more than 30 seconds old. Give the job and the monitor their own users, never the replicator's. The signal outlivesa's monitoring; its error is the clock offset between the regions.
Clocks¶
No time arithmetic crosses the clusters in the data path: lag_seconds compares a's commit time with the clock of the same a node. A message's timestamp on b is when it arrived on b, so after an outage hours of records arrive within minutes of timestamps, and a delay child on b delays from arrival. Certificate validity is checked against a's clock. Keep both clusters on NTP or the cloud's time service.
Next steps¶
- Fail over to the recovery cluster: when
a's region is lost. - Fail back to the original cluster: when it is back.
- Manage remotes: credentials, rotation and the security model.