System designby Learnastra

Concept lesson · Foundations

Distributed systems: scalability, reliability, availability and efficiency

By Anup Rai

Start here

Definition

A distributed system consists of independent computers that coordinate by exchanging messages. Its quality must be assessed separately: scalability concerns increased workload, reliability concerns correct service over time, and availability concerns whether service is usable when requested. Efficiency measures useful work per resource spent; manageability concerns safe diagnosis, repair and change.

Why it matters: Running on several computers introduces partial failures: the application can be alive while the database is unreachable. Separate quality targets tell you which failure matters and how to respond.

The visual modelScalability, reliability, availability, efficiency, and operability

Scalability, reliability, availability, efficiency and manageability are distinct quality attributes. Each needs its own definition and measurement.

Scalability, reliability, availability, efficiency, and operabilityScalability, reliability, availability, efficiency and manageability are distinct quality attributes. Each needs its own definition and measurement. Scaling asks whether 500 checkout requests/s can become 2,000 while maintaining latency. Reliability asks whether one intended purchase yields O17 and one correct charge, including retries. Request availability counts successful eligible checkouts; one million attempts at 99.9 percent permits 1,000 unsuccessful attempts. Efficiency measures useful work per resource; manageability covers diagnosing, repairing and changing service safely.Checkout service: five measurable quality attributesSCALABILITY500 to 2,000 requests/sMaintain the agreed latency at higher loadRELIABILITYOne purchase becomes O17No duplicate order or charge after a retryAVAILABILITY99.9% of eligible attempts1,000 misses per 1,000,000 attemptsEFFICIENCYUseful work per resourceMeasure CPU, network and cost per checkoutMANAGEABILITYOperate and change safelyDetect, repair, roll back and verifyCorrectness and availability depend on the defined user outcome, not just a responding server.
Read the diagram step by step
  1. Scaling asks whether 500 checkout requests/s can become 2,000 while maintaining latency.
  2. Reliability asks whether one intended purchase yields O17 and one correct charge, including retries.
  3. Request availability counts successful eligible checkouts; one million attempts at 99.9 percent permits 1,000 unsuccessful attempts.
  4. Efficiency measures useful work per resource; manageability covers diagnosing, repairing and changing service safely.

Worked example

Order O17 is a $25 purchase. A second application server can accept traffic after the first fails, but a repeated request still needs to recover O17 rather than create a second $25 charge.

Key takeaways

  • Reachable processes do not prove a correct user outcome.
  • More application servers do not remove a shared database bottleneck.
  • State the failure being tolerated and the capacity left afterward.

You will learn to

  • Explain each system quality using an observable user outcome.
  • Calculate an availability/error budget and surviving capacity.
  • Identify why adding machines can leave a bottleneck unchanged.

Practice in this chapter

8 interview questions with model answers and follow-ups.

Go to interview practice

Useful foundations: Capacity estimation: throughput, latency, concurrency and storage

Workload and timing examples are interview assumptions.

01Distributed system quality attributes: definitions

A distributed system consists of independent computers that coordinate by exchanging messages. One part can fail while others keep running; this is a partial failure. An online shop might run its request-handling application on two machines, keep live copies of orders on several database machines, and send receipts through a background worker. The customer sees one checkout experience, even though these parts can fail or respond at different times.

The qualities below answer different questions about that experience. Scalability asks whether the service can handle a larger workload while maintaining its targets. Reliability is the ability to perform the specified function correctly under stated conditions over a period of time. Availability is the degree to which the service is usable when requested, often measured as successful eligible requests divided by all eligible requests. Efficiency asks how much useful work it gets from its resources. Manageability asks how safely operators can observe, configure, operate, and change it. The related term serviceability focuses on diagnosing and repairing faults.

Quality Question for checkout Example design decision Remaining limit
Scalability Can 500 requests/s become 2,000 at the same latency? Distribute independent application work One hot inventory row may remain serial
Reliability Does one purchase create the intended order and charge? Save request identity and reconcile payment outcomes A remote payment requires its own retry contract
Availability Can an eligible customer complete checkout now? Keep spare replicas and fail over safely A partition may require refusal rather than unsafe writes
Efficiency How much CPU, data and communication does one order consume? Remove repeated lookups and batch safe work Larger batches may increase wait time
Manageability Can an operator diagnose and repair O17? Trace IDs, durable states and staged rollouts Automation still needs safe thresholds

02Vertical scaling, horizontal scaling and serial bottlenecks

At first, one application server handles 500 checkout requests/s. Vertical scaling replaces it with a larger machine: more CPU, memory, or faster disks. It can be the simplest improvement, but hardware has practical limits and a single machine still fails as one unit.

Concept in focusBigger machine or more machines?

Machine size represents resources per instance; separate boxes represent independent instances. Neither change removes a shared database bottleneck.

Bigger machine or more machines?Machine size represents resources per instance; separate boxes represent independent instances. Neither change removes a shared database bottleneck. Compare one enlarged instance with work spread across three instances. Vertical scaling replaces a two-CPU instance with an eight-CPU instance. Horizontal scaling routes work across three two-CPU instances.Vertical: enlarge one machine2 CPU8 CPUHorizontal: distribute work across machinesRouter2 CPU2 CPU2 CPUThree independent instances

Remember: Vertical changes the size; horizontal changes the count.

Read the diagram
  1. Compare one enlarged instance with work spread across three instances.
  2. Vertical scaling replaces a two-CPU instance with an eight-CPU instance.
  3. Horizontal scaling routes work across three two-CPU instances.
Try from memoryWhich approach spreads work across several server instances?

Horizontal scaling adds independent instances. It can tolerate an instance loss only if routing, surviving capacity and state management support it.

Horizontal scaling adds machines. Put two application servers behind a load balancer, and either can handle a request if essential state is stored outside the process. This can grow application capacity and tolerate one application failure if the survivor can meet the admitted workload. It does not automatically double database write capacity.

Some work remains serialized: operations must take turns because they update the same protected state. Adding application machines does not remove that ordering requirement. This matters both for the time one checkout takes and for how many checkouts can update the same inventory record.

For a separate latency calculation, suppose one request spends 80 ms on parallelizable work and 20 ms executing a serialized operation on one inventory key, excluding queue wait. Making the first part four times faster yields 80/4 + 20 = 40 ms, a 2.5× improvement, not 4×. Even infinitely fast application work cannot eliminate the remaining 20 ms. This is the intuition behind a serial bottleneck: improve the part that limits the actual operation.

The workload also matters. Adding nodes can help independent product lookups while thousands of purchases of the same final item still contend on one record. Measure distribution, not just total QPS.

03Availability and error-budget calculations

An error budget is the amount of unsuccessful service allowed by the chosen availability target over a defined measurement window. The target supplies the permitted fraction; the number of requests or the duration of the window turns it into a count or time allowance. Choose that denominator before interpreting an outage.

At 10:00 the only order database stops responding. Automated detection fires at 10:01. An operator finishes failover and verifies writes at 10:07. Checkout was unavailable for seven minutes, not merely the six minutes spent repairing after detection. Monitoring delay is part of the user impact.

For a simple recurring up/down model, availability can be approximated by mean uptime / (mean uptime + mean downtime). Real services have partial and correlated failures, so a single formula is not a substitute for measuring user requests. Faster detection and repair can improve availability even when the underlying failure frequency is unchanged.

Dependencies also affect the result. In a deliberately simplified model, if two required dependencies are independently available 99.9% of the time, the path is available 0.999 × 0.999 = 99.8001% of the time before other failure sources. Redundant alternatives instead help only when at least one is usable and routing can reach it. Shared power, bad configuration and overload make failures correlated, so multiplying advertised service percentages is not a production reliability proof.

Worked example diagramTwo application servers still converge on one inventory writer. The shared write can limit scalability even while the application tier has spare CPU.
Distributed systems: scalability, reliability, availability and efficiency: architecture diagram1. Purchase request O17 to 2. Load balancer: checkout request; 2. Load balancer to 3. Application A: healthy instance; 2. Load balancer to 4. Application B: another healthy instance; 3. Application A to 5. Shared inventory writer: reserve inventory; 4. Application B to 5. Shared inventory writer: same shared writer; 5. Shared inventory writer to 6. Receipt worker: receipt after committed order1 → 2: checkout request2 → 3: healthy instance2 → 4: another healthy instance3 → 5: reserve inventory4 → 5: same shared writer5 → 6: receipt after committed order01Purchase request O1702Load balancer03Application A04Application B05Shared inventorywriter06Receipt worker
  1. 1 → 2checkout requestPurchase request O17 → Load balancer
  2. 2 → 3healthy instanceLoad balancer → Application A
  3. 2 → 4another healthy instanceLoad balancer → Application B
  4. 3 → 5reserve inventoryApplication A → Shared inventory writer
  5. 4 → 5same shared writerApplication B → Shared inventory writer
  6. 5 → 6receipt after committed orderShared inventory writer → Receipt worker

04Reliability, durability and failure domains

If order O17 commits but its response is lost, a retry can create O18 and charge again. Save the result under a stable request ID so a retry returns O17. Save the order and request result in one atomic database transaction: both commit or neither does. An external payment is outside that transaction. Reuse the same payment identifier under the provider’s retry rules, and check an uncertain result before issuing another charge.

Durability is retention of acknowledged data. Replicated records can survive a machine loss if the acknowledgment and recovery protocol make that promise. A backup can restore an earlier state after accidental deletion. Both require verification; merely drawing duplicate cylinders does not prove an acknowledged purchase survives.

A failure domain is a set of components that one event can disable together, such as machines sharing a power supply or deployment zone. A network partition prevents some machines from communicating even though they may still be running. Replica placement must match the failures the service is meant to survive.

The following failure sequence shows which records and identifiers must survive a lost response. (1) Purchase key K17 requests $25 and the order authority records O17. (2) Payment action charge-O17 produces confirmed provider charge C81. (3) The application response is lost. (4) Retrying K17 returns O17/C81 rather than allocating O18 or a new charge identity. If the provider response was lost instead, the charge remains unknown until lookup or the provider’s documented same-key retry resolves it. A timeout establishes uncertainty, not failure.

05Resource efficiency and communication cost

Efficiency is useful outcomes divided by the resources spent. For O17, ten internal RPCs—remote procedure calls—may each transfer a small record. One giant catalog transfer may use fewer messages but far more bytes. Count both messages and data size, then account for network distance and repeated work.

Suppose design A makes ten sequential 5 ms calls and design B makes two 20 ms calls. Their network wait contributions are about 50 and 40 ms respectively in this simplified example. A third design could batch data into one call, but might waste bytes or postpone the response. Message count alone does not identify the best design.

Symptom Likely resource to inspect Example improvement
CPU saturated on every server Computation per request Remove repeated parsing or cache a safe result
Database reads dominate Query plan and indexes Fetch O17 by an indexed identifier
Large transfers dominate Bytes and distance Compress or deliver static bytes nearer readers
Only one partition is hot Work distribution Revisit ownership, split a hot workload

Mixed machine sizes, topology, and uneven load make ideal linear speedup unlikely. Compare designs with the same workload and objective.

06Manageability, monitoring and safe change

An operator should be able to answer what failed, which customers are affected, and which action is safe. Attach one request/trace ID to O17 across services, record state transitions without payment secrets, and measure both successful outcomes and latency. A health endpoint that only says the process is alive does not prove orders can commit.

Use different controls for different problems. Readiness decides whether an instance receives new requests. A restart policy decides when to restart its process. Admission control limits accepted work so existing requests can finish. If a dependency fails, accept less work where necessary; restarting otherwise healthy application processes will not fix that dependency.

Roll out a new version to a small fraction first, compare outcomes, and retain a rollback path. A database change should let old and new application versions coexist during the rollout. Stop assigning new work to a known dead instance, and bound or shed the excess traffic if survivors lack capacity. Separately, avoid ejecting or repeatedly restarting every live instance merely because a shared dependency is slow: that reaction can reduce useful capacity further. Readiness, restart policy and admission control have different jobs.

Candidate explanation: “I separate checkout availability from order correctness. I can temporarily refuse new purchases when I cannot confirm which database node is allowed to update inventory, while keeping browsing available. I add application redundancy, make retries return the original order, and measure the full checkout outcome. My recovery plan includes detection, failover, validation, and enough remaining capacity.”

Practise the interview questions

Say your answer aloud before opening the model answer. Then answer the follow-up and compare the reasoning.

Foundation · Question 1

What is the difference between reliability and availability?

Reveal a model answer

“Availability asks whether an eligible checkout operation can complete under its success definition. Reliability asks whether the service performs its specified function correctly over time and under promised conditions. A reachable system that double-charges an order is incorrect; that purchase must also count as unsuccessful in an end-to-end availability measure. I define the outcome and measurement window rather than treating reachability as either guarantee.”

What the answer must demonstrate: Use the same example for both qualities.

Foundation · Question 2

When would you scale vertically before sharding?

Reveal a model answer

“If the database fits on one larger instance and measured CPU, memory, or I/O is the bottleneck, vertical scaling can buy capacity with a smaller operational change. I would also keep redundancy and test the new capacity. I shard when independent data needs to exceed that practical limit.”

What the answer must demonstrate: Separate physical resources from contention.

Applied · Question 3

Why does doubling application servers not double checkout throughput?

Reveal a model answer

“They may still share the same database, lock, or downstream service. I trace a purchase and measure where time and work accumulate. Adding application capacity helps only the work those instances own; the shared inventory writer may remain the limiting resource.”

What the answer must demonstrate: Find the shared bottleneck.

Applied · Question 4

What does 99.9% availability permit?

Reveal a model answer

“First I would define the measure. Over a 30-day time-based window, 0.1% is 43.2 minutes. Over a million eligible requests, it is 1,000 unsuccessful attempts. These budgets are not interchangeable when traffic changes through the day.”

What the answer must demonstrate: Define eligible and successful requests.

Applied · Question 5

Why include detection time in a recovery plan?

Reveal a model answer

“The customer experiences the outage before the operator starts repairing. If detection takes one minute and verified failover takes six more, checkout is unavailable for seven. I improve both detection and repair and practise the complete sequence.”

What the answer must demonstrate: Measure end-to-end recovery.

Applied · Question 6

Do two copies guarantee durability?

Reveal a model answer

“No. I need to specify when a write is acknowledged, whether the second copy is durable, and which failures it survives. Copies in the same failure domain may disappear together, and a bad deletion can replicate to both. I also need backups and tested recovery.”

What the answer must demonstrate: Name the failure being tolerated.

Applied · Question 7

Is fewer network messages always more efficient?

Reveal a model answer

“No. One message may contain a huge unused payload, while several small messages may run in parallel. I compare bytes, round trips, CPU, and end-to-end latency for the same user operation. Reducing repeated calls can help, but the workload decides.”

What the answer must demonstrate: Count bytes and sequential waits, not just arrows.

Applied · Question 8

What makes a system manageable in an interview answer?

Reveal a model answer

“I show how an operator diagnoses one failed order using a trace identifier and durable states, how alerts reflect failed purchases, and how a rollout can be stopped or reversed. I include schema compatibility and verify recovery rather than ending the design at deployment.”

What the answer must demonstrate: Explain a concrete operator action.

Blank-page exercise · 15 minutes

Build the answer yourself

Explain why a reachable checkout can be unreliable, then redesign it to survive one application failure.

  • Give one example for each of the five qualities.
  • Calculate a stated availability budget.
  • Trace a lost-response retry for one purchase.
  • Identify one shared failure domain and one serial bottleneck.

Check that each component and design decision follows from your requirements and workload.

Recall the key ideas

Answer from memory before opening each card. Explain why the choice works and what it costs. Revisit missed cards tomorrow.

Distributed systems: scalability, reliability, availability and efficiencyFive system qualitiesRecall first, then reveal

Scalability: more work; reliability: correct work; availability: usable now; efficiency: resource cost; manageability: safe diagnosis and change.

Ask five different questions.

Return to lesson
Distributed systems: scalability, reliability, availability and efficiencyAvailability budgetRecall first, then reveal

Choose a request-based or time-based definition and a window before calculating.

Define the denominator.

Return to lesson
Distributed systems: scalability, reliability, availability and efficiencyWhy might adding application servers fail to speed up checkout?Recall first, then reveal

If every server still waits on the same overloaded database, adding servers leaves the bottleneck in place. Distribute or reduce the limiting work.

Find the bottleneck before adding machines.

Return to lesson

Final revision

Summary and interview notes

A distributed service must be evaluated at the user-visible operation, not by counting reachable machines. Scalability, reliability, availability, efficiency and manageability describe different qualities, and each needs its own workload, failure model and measurement.

Remember these points

  • Adding machines helps work that can run independently; updates to one heavily used key may still have to run one at a time.
  • A request-based availability budget differs from a time-based outage budget.
  • Reliable retries reuse the original operation ID and stored result. If an external action such as a charge has an unknown outcome, check its status before attempting a new action.
  • Copies protect only against the failures covered by their placement, acknowledgment and recovery protocol.

Interview tips

  • Use one operation to contrast the five qualities, then explain how each is measured.
  • Separate a per-request latency speedup from aggregate throughput and hot-key capacity.
  • Include detection, failover and verified service recovery in the outage timeline.

Important qualifications

  • End-to-end success should count incorrect results as failures; process reachability alone is a weaker metric.
  • Independence-based availability arithmetic is a simplified model; shared dependencies and correlated failures require direct measurement.
  • Stop routing to failed nodes. Limit accepted work so redirected traffic does not overload the survivors.

Technical references

Practice marks stay in this browser.