Concept lesson · Foundations
Distributed systems: scalability, reliability, availability and efficiency
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.
Scalability, reliability, availability, efficiency and manageability are distinct quality attributes. Each needs its own definition and measurement.
Read the diagram step by step
- 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.
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 practiceUseful foundations: Capacity estimation: throughput, latency, concurrency and storage
Workload and timing examples are interview assumptions.
Dotted concept links open the relevant explanation in a new tab.
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.
Machine size represents resources per instance; separate boxes represent independent instances. Neither change removes a shared database bottleneck.
Remember: Vertical changes the size; horizontal changes the count.
Read the diagram
- 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.
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.
- 1 → 2checkout requestPurchase request O17 → Load balancer
- 2 → 3healthy instanceLoad balancer → Application A
- 2 → 4another healthy instanceLoad balancer → Application B
- 3 → 5reserve inventoryApplication A → Shared inventory writer
- 4 → 5same shared writerApplication B → Shared inventory writer
- 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.
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.”
Interviewer follow-up
Can you preserve correctness while losing availability?
Reveal the follow-up answer
Yes. Refusing a purchase when the service cannot confirm which node may update inventory avoids accepting an order it cannot safely reserve stock for, but the customer still cannot complete the operation.
What the answer must demonstrate: Use the same example for both qualities.
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.”
Interviewer follow-up
What does vertical scaling not fix?
Reveal the follow-up answer
A logical lock on one hot inventory record can remain a serial bottleneck. Bigger hardware is not a concurrency protocol.
What the answer must demonstrate: Separate physical resources from contention.
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.”
Interviewer follow-up
What if browsing scales but checkout does not?
Reveal the follow-up answer
That is plausible because browsing can distribute read work while checkout changes shared inventory. I would size and design those operations separately.
What the answer must demonstrate: Find the shared bottleneck.
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.”
Interviewer follow-up
Can a fast error count as successful?
Reveal the follow-up answer
Only if it is a valid business response under the defined metric, not because the network responded quickly. An infrastructure refusal of a valid purchase is an unavailable outcome.
What the answer must demonstrate: Define eligible and successful requests.
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.”
Interviewer follow-up
Would aggressive health checks always help?
Reveal the follow-up answer
No. Noisy checks can eject or restart live capacity during a shared dependency incident. Stop sending work to a known dead instance, but use bounded admission and careful failure thresholds to prevent overload from cascading through the survivors.
What the answer must demonstrate: Measure end-to-end recovery.
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.”
Interviewer follow-up
What is a failure domain?
Reveal the follow-up answer
A set of resources that can fail together because they share a dependency, such as power, a rack, a zone, or an administrative change.
What the answer must demonstrate: Name the failure being tolerated.
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.”
Interviewer follow-up
What changes across regions?
Reveal the follow-up answer
Each sequential round trip can cost substantially more time because of distance. I would reduce cross-region dependencies on the critical path and measure the actual network.
What the answer must demonstrate: Count bytes and sequential waits, not just arrows.
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.”
Interviewer follow-up
Which metric would you alert on first?
Reveal the follow-up answer
The user-facing purchase-success or latency objective, supported by component metrics to locate the cause. A low-level CPU signal alone does not establish customer impact.
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 lessonDistributed 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 lessonDistributed 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 lessonFinal 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
- Google SRE service-level objectivesUser-facing measures of availability, latency, and correctness.
- Google SRE addressing cascading failuresFailure amplification, remaining capacity, and recovery behavior.
Practice marks stay in this browser.