Replication Flow Control and Backpressure
LESSON
Replication Flow Control and Backpressure
By the end of this lesson, you will be able to...
distinguish send, durable, apply, and retention backlog in a replicated shard;
trace how a leader should protect a quorum-critical follower while another replica recovers;
choose when to cap catch-up traffic, send a snapshot, or slow foreground writes.
Idea in one sentence: Backpressure keeps a replica problem bounded by deciding whose backlog may grow, how far it may grow, and when clients must slow down.
Core Insight
At market open, md-db-2 leads Harbor Point's shard 184. Its local follower, md-db-4, is almost current. Its remote follower, ny-db-3, is recovering over a congested link. Reservation writes suddenly rise.
The easy response is to let the leader queue every new log entry for New York and send as much as the network accepts. For a moment, the API looks fast. The cost is hidden in leader memory, retained log segments, and crowded replication workers. Eventually, even healthy quorum heartbeats compete with bulk catch-up traffic.
The better question is not “how do we make every follower equally fast?” It is “which path proves the service contract right now, and how do we keep an optional recovery path from consuming that proof?”
For this shard, commits need the leader and a local quorum follower. The remote replica improves recovery and failback, but it is not on the normal commit path. Backpressure protects the majority path first, then gives the remote replica a bounded way to recover. If the majority path starts losing ground, backpressure must reach clients before invisible queues become a failover or storage crisis.
The Situation: One Word, Four Backlogs
“Replication lag” is too vague to operate. A log entry moves through several stages for each follower:
leader last index
-> sent to follower
-> durably stored by follower
-> applied by follower's state machine
For ny-db-3, the leader might expose these illustrative positions during the burst:
| Position or resource | Value | What it tells us |
|---|---|---|
| Leader last index | 813,620 | Latest entry accepted locally. |
| Sent index | 811,900 | Data waiting in the network or follower receive path. |
| Durable index | 809,940 | Last entry the follower can recover after its own crash. |
| Applied index | 809,100 | Last entry visible to follower-served reads. |
| Retained log floor | 809,940 | Oldest point still needed for incremental catch-up. |
Each gap changes the diagnosis.
- A large gap from leader to sent points to sender scheduling, congestion, or a full in-flight window.
- A large gap from sent to durable points to follower receive or flush capacity.
- A large gap from durable to applied can leave commit safety intact while making follower reads stale.
- A very old retained-log floor consumes disk on the leader because the lagging follower still needs history.
The numbers are a teaching trace, not a measurement. Their job is to show why one “lag seconds” graph cannot decide whether to tune a network window, protect disk, stop serving follower reads, or prepare a snapshot.
The Initial Model: Send Everything Faster
The initial model is reasonable: if New York is behind, give it more sender buffers and larger batches. This works when the follower has a short, known delay and enough disk and network capacity to catch up. A larger batch can reduce per-message overhead.
It breaks when the follower is slower than the incoming write stream. A leader that can append 18,000 entries per second but can deliver only 7,000 net entries per second to New York creates a growing deficit every second. More buffering changes where that deficit lives; it does not remove it.
write stream: +18,000 entries/s
remote catch-up: -7,000 entries/s
unbounded deficit: +11,000 entries/s
Without a bound, the system accumulates unsent bytes in memory and old log segments on disk. The remote follower may eventually need a snapshot, but now the snapshot is larger and the leader has spent scarce resources postponing that decision. Worse, a shared queue can delay md-db-4 heartbeats and acknowledgements, turning an optional follower's problem into a quorum problem.
The missing mechanism is flow control: the leader gives each follower a limited amount of data it may have in flight, replenishing that credit only as the follower acknowledges durable progress. Backpressure is the policy above that mechanism: it decides which credits shrink, which traffic has priority, and when write admission changes.
The Moving Parts and Their Priorities
Harbor Point separates replication work into three classes:
| Class | Example | Priority | Why |
| --- | --- | --- |
| Control | term changes, heartbeats, acknowledgements | Highest | They establish current authority and quorum health. |
| Quorum-critical data | appends to md-db-4 | High | They determine normal commit latency and durable majority progress. |
| Recovery data | bulk catch-up to ny-db-3 | Bounded | It improves resilience but must not crowd out the first two classes. |
This is a situated choice, not a rule that remote replicas never matter. A design whose acknowledgement contract includes the remote region must treat that replica as critical. The point is to classify a follower by the guarantee it currently supports, not by its geographic label.
For the example, the leader starts with a 4 MB in-flight window for md-db-4 and a 512 KB window for recovering ny-db-3. Credits are illustrative. They keep the leader from letting one slow receiver reserve unlimited memory and worker time. The leader also reserves scheduling capacity for heartbeats and quorum appends; a full recovery window may delay its own bulk entries, but it may not starve the control lane.
A Worked Control Loop
At 09:30:00, writes rise and New York falls behind. The leader evaluates this small policy every few seconds:
If a quorum follower's durable lag or acknowledgement delay grows:
slow write admission and preserve control/quorum lanes.
Else if an optional follower's retained-log cost is too high:
stop incremental catch-up and re-seed it with a snapshot.
Else if an optional follower's window is full:
pause its new sends until it acknowledges more durable progress.
Else:
continue normal replication.
Here is one trace with illustrative thresholds.
| Time | md-db-4 durable lag |
ny-db-3 state |
Leader decision | Reason |
|---|---|---|---|---|
| 09:30:00 | 8 ms | Window full; remote durable lag grows | Cap New York at 512 KB. | Majority path is healthy; do not add remote queue debt. |
| 09:31:30 | 11 ms | Retention due to New York reaches 18 GB | Continue bounded catch-up. | Still inside the 20 GB policy budget. |
| 09:32:10 | 14 ms | Retention reaches 20 GB; catch-up rate stays below write rate | Stop incremental catch-up; schedule snapshot. | More retained log will not close the deficit cheaply. |
| 09:33:00 | 55 ms | Snapshot work is isolated | Throttle foreground writes. | The quorum path, not the optional path, now lacks margin. |
The order matters. At 09:30, stopping all writes would throw away availability without protecting a guarantee that is actually at risk. At 09:33, continuing normally would let the quorum follower fall farther behind, which threatens commit latency and the timely authority evidence used by lease-read fast paths. The controller therefore moves the pain to the smallest safe boundary first: remote flow, then recovery mode, then client admission.
So far: flow control keeps recovery traffic finite; backpressure decides who yields when finite resources are contested. Snapshotting is not punishment for a slow follower. It replaces an ever-growing incremental repair job with a bounded recovery operation.
What the User Sees—and What the System Knows
The user may see increasing write latency, a retryable overload response, or a temporarily stale regional read. Those outcomes differ, and the system should not hide them behind a generic “database slow” alert.
| Signal | Meaning | First action |
|---|---|---|
| Quorum durable lag grows | The normal acknowledgement path is losing margin. | Protect control traffic; reduce write admission. |
| Optional follower in-flight bytes stay full | The receiver cannot absorb its assigned recovery rate. | Hold its credit; inspect network and follower I/O. |
| Retained log bytes grow toward a limit | Incremental history is becoming expensive to keep. | Prepare snapshot or detach/reseed under the membership policy. |
| Applied lag grows while durable lag is small | Follower reads may be stale, even if its log is durable. | Remove it from freshness-sensitive reads or require a version check. |
| Heartbeat delay rises with bulk traffic | Queues are mixing control and recovery work. | Reserve or isolate the control lane. |
This evidence also avoids a common wrong fix: increasing every buffer after a latency spike. Larger queues can smooth a short burst. They cannot manufacture follower flush bandwidth, disk space, or WAN capacity. When the sustained input rate exceeds the recovery rate, bigger queues lengthen the time before the system admits that it needs a different response.
Cost, Limits, and Signals
Backpressure improves safety and recovery posture by bounding debt. The trade-off is normal-path throughput and earlier visible overload. A snapshot costs transfer work and time before the follower becomes useful again. A detached follower reduces redundant copies until it catches up. These are design choices under a stated durability and availability contract, not automatic fixes.
The boundary is the recovery story. If the leader loses retained history before the follower can receive it, a snapshot must come from a known committed state. If the snapshot source or validation path is missing, a per-follower window only delays an unrecoverable condition. Membership status matters too: do not casually remove a lagger that is still required by the committed quorum configuration. Use the reconfiguration procedure from the next lesson when the replica set itself must change.
The key signals are not decorative capacity metrics. They answer operational decisions: which stage is slow, whether a follower is still eligible for a read, how much retention is being consumed, and whether the quorum path needs client-facing pressure now.
Check Your Understanding
Check: A follower's durable index nearly matches the leader, but its applied index is far behind. Is the first problem necessarily commit durability?
Think first, then reveal.
Answer: No. The follower may have durably stored the log while applying it slowly. This can make its reads stale, but it is different from a missing durable acknowledgement. Inspect the serve-read contract before treating it as a quorum-commit failure.
Check: An optional remote follower falls farther behind every second, while the quorum follower remains healthy. Should the leader immediately reject all writes?
Think first, then reveal.
Answer: Not necessarily. First cap the optional stream and bound retained history; switch to snapshot recovery when incremental catch-up cannot win. Foreground throttling becomes necessary when the quorum path or a stated recovery guarantee is at risk.
Practice: Set a Bounded Policy
For a three-replica service, write a four-line backpressure policy. Name one quorum-critical follower, one optional recovering follower, a retained-log threshold, and the signal that begins write throttling.
A good answer should mention: separate priority for heartbeats and critical appends; a finite in-flight window for recovery; a snapshot or reseed decision when retained history crosses a limit; and a client-facing throttle tied to growing quorum durable lag, not merely any follower lag.
Connections
- Clocks, Leases, and Safe Reads shows why delayed quorum communication can close a lease-read fast path before any safety violation occurs.
- Membership Changes and Replica Set Evolution explains why a lagging member remains part of the quorum until a committed reconfiguration changes that fact.
- Backup, Snapshots, and Recovery Semantics provides the recovery boundary when incremental log catch-up is no longer economical or possible.
Resources
- [PAPER] In Search of an Understandable Consensus Algorithm — Focus: Inspect per-follower
nextIndex,matchIndex, append progress, and snapshot installation as ingredients of catch-up control. - [DOC] PostgreSQL Replication Configuration — Focus: Connect replication slots and sender configuration to retained-WAL pressure caused by lagging standbys.
- [DOC] etcd Tuning — Focus: Relate heartbeat and election timing to the importance of preserving timely quorum communication under load.
Key Takeaways
- Sender, durable, apply, and retention gaps are different backlogs and need different interventions.
- Protect control traffic and quorum-critical replication before spending unbounded resources on a recovering optional follower.
- Backpressure makes overload visible at a chosen boundary—follower credit, snapshot recovery, or client admission—instead of letting it surface later as disk exhaustion or fragile failover.