Learning pathsA
Patterns

Scaling Writes

A write path is limited by the slowest sustained stage, not by the length of its queue.

A write path is limited by the slowest sustained stage, not by the length of its queue. Scaling it requires understanding write amplification, partition balance, durability, and overload behavior. A queue can smooth bursts and decouple callers, but it cannot make a database or external provider process more work per second indefinitely.

Learning goals

Write amplification; batching; append logs; partitioning; hot keys; queues; backpressure; admission control; bounded backlog; replay.

The mechanism at a glance

Producers → Admission budget (incoming rate); Admission budget → Durable log (accepted work); Durable log → Partition workers (partitioned reads); Partition workers → Storage owners (bounded writes); Durable log → Lag monitor (oldest age); Lag monitor → Admission budget (throttle)
Scroll to inspect the diagram, or open it at full size.

Figure — Producers → Admission budget (incoming rate); Admission budget → Durable log (accepted work); Durable log → Partition workers (partitioned reads); Partition workers → Storage owners (bounded writes); Durable log → Lag monitor (oldest age); Lag monitor → Admission budget (throttle)

The numbered components identify responsibilities. Follow the labeled arrows rather than treating the numbers as a global execution order. The scenario later in this lesson shows one concrete sequence.

Step-by-step reasoning

1. Count the real write work

One incoming event may update several indexes, replicas, counters, and projections. Measure bytes and storage operations, not just API requests. Remove unnecessary synchronous work and use append-oriented records when the access pattern permits. Batching amortizes overhead but adds waiting time and makes failure units larger.

2. Acknowledge at a defined boundary

If the API returns success before durable acceptance, a crash can lose accepted work. Persist to a durable log or transactionally store intent before acknowledging according to the stated contract. Consumers then apply derived work asynchronously. Define retention and replay so a downstream outage does not outlast recoverable history.

3. Partition for independent progress

Choose a key that preserves only the ordering actually needed. Ordering every event globally unnecessarily serializes the system. Per-account or per-campaign ordering may be sufficient. Watch skew: a single popular campaign can overload one partition even with hundreds of idle partitions elsewhere.

4. Apply backpressure before collapse

Measure oldest-message age, not just queue depth. Bound producers, consumer concurrency, and retries. Reject, defer, sample, or prioritize work according to its importance. Billing events may require durable rejection and retry; optional telemetry may permit sampling. Make this a product contract rather than a silent data-loss policy.

Contracts and state

The following sketch makes the decision boundary concrete. Field names and capacity assumptions are illustrative; adapt them to the stated product contract.

Contract / pseudocode
Arrival = 12,000 events/s
Sustainable service = 8,000 events/s
Backlog growth = 4,000 events/s
After 1 hour = 14.4 million queued events
Scaling workers helps only if the bottleneck can process more.

Worked example

A telemetry service has short bursts at 12,000/s but normally receives 3,000/s. An 8,000/s consumer tier can absorb a bounded burst and drain later. If 12,000/s lasts all day, the same design is unstable. Estimate burst duration and available retention. Add capacity at the actual bottleneck, reduce work per event, or throttle admission; do not label the queue “infinite.”

At a sustained 4,000 events/s deficit, backlog reaches 14.4 million in one hour.
Scroll to inspect the diagram, or open it at full size.

Figure — At a sustained 4,000 events/s deficit, backlog reaches 14.4 million in one hour.

Failure walkthrough

Retries during a storage slowdown multiply write pressure. Use exponential backoff with jitter, per-tenant limits, and a bounded retry budget. Poison events need an inspectable quarantine path so they do not block a partition forever. Replay must preserve original event identities; generating new IDs turns a repair into duplicate business activity.

Traffic exceeds service → Queue age increases → Admission limit engages → Useful work remains stable → Spare capacity drains backlog
Scroll to inspect the diagram, or open it at full size.

Figure — Traffic exceeds service → Queue age increases → Admission limit engages → Useful work remains stable → Spare capacity drains backlog

Decisions and trade-offs

TechniqueGainCost
BatchingLess overhead per eventLatency and larger retry units
PartitioningParallel ownersCross-key coordination
Async projectionShorter acceptance pathLag and replay operations
Admission controlProtects stable capacityExplicit rejection or delay

Check your understanding

An operator adds ten consumers but throughput stays fixed. What should they inspect next?

Show answer and explanation

Answer: Inspect partition parallelism, key skew, database write capacity, provider quotas, and lock contention. Extra consumers cannot exceed a fixed downstream budget or split one ordered hot partition automatically.

Transfer to a new scenario

Handle a burst of telemetry without claiming the queue increases downstream capacity.

At 12k incoming and 8k processed events/s, what happens during a sustained hour?

Continue the connection

Study Kafka and explain which guarantee from this lesson carries into that topic.

Distinguish admission from completion

A queue lets the API acknowledge durable acceptance before expensive processing, but the durability claim depends on what was committed before the acknowledgement. Store a job plus outbox or append to a durable log under an explicit contract. The worker pool then consumes within downstream concurrency and throughput limits.

Batching amortizes fixed overhead but trades added wait time and larger retry units. Flush on item count, byte count and time, and keep a stable batch or item identity. Partition by the ordering key; random partitioning spreads load but loses per-entity order unless another layer restores it.

Monitor oldest-event age alongside throughput. A queue can look stable in length while one low-priority tenant starves. Bound per-tenant buffers and apply a documented fairness policy. When sustained arrival exceeds service, choose backpressure, rejection, delayed completion or reduced work; buffering alone cannot settle the deficit.

A decision worksheet for Scaling Writes: read the mechanism and its guarantee together.
Scroll to inspect the diagram, or open it at full size.

Figure — A decision worksheet for Scaling Writes: read the mechanism and its guarantee together.

Operational sketch

Contract / pseudocode
accept -> durable log -> partition workers -> sink
flush if count >= 100 OR bytes >= limit OR age >= 20ms
commit sink effect before safe acknowledgement
retry using original event identity

A tempting mistake

Increasing workers against an exhausted connection pool can lower throughput through contention and timeout retries. Scale after identifying the bottleneck and reserve capacity for recovery traffic.

Transfer exercise

Why might larger batches worsen p99 latency?

Show answer and explanation

Answer: Items wait for a batch to fill and one slow item or retry can delay a larger group. Use a time cap and measure both throughput and end-to-end age.

Your study notes