0%
Data-Intensive Applications
Distributed Data
Encoding and Evolution
Batch Processing
Stream Processing
Data Quality and Governance
Operational Patterns
Reliability, Scalability, Maintainability
A system is reliable when it continues to work correctly even when things go wrong. Not "nothing ever fails" but rather "failures happen and the system handles them gracefully." This distinction matters because at scale, failure is not an exception. It is a statistical certainty. A cluster of 10,000 hard drives will see multiple disk failures every single day. The question is never "will something break?" but "when it breaks, does the user notice?"
The term to internalize is fault tolerance, not fault prevention. A fault is one component deviating from its spec (a disk dying, a network packet dropping, a process crashing). A failure is when the system as a whole stops providing the service the user expects. Reliable systems convert faults into non-events by designing redundancy and recovery into every layer.
Hardware Faults
Hardware faults are the easiest to reason about because they follow well-understood probability distributions. A hard drive has a mean time to failure (MTTF) of about 10-50 years. On a cluster with 10,000 disks, you should expect at least one disk failure per day on average. Power supplies fail. RAM develops bit errors. Network cables get unplugged by someone tripping over them in the datacenter.
The standard response to hardware faults is redundancy. RAID arrays tolerate disk failures by mirroring data across multiple drives. Dual power supplies with separate electrical feeds survive a power event on one circuit. Hot-spare machines sit idle waiting to take over when an active node goes down. ECC memory detects and corrects single-bit errors before they corrupt data. None of this is exotic. These are table-stakes techniques that keep systems running through individual component failures.
What changed in the last decade is the shift from hardware redundancy to software-level fault tolerance. Cloud platforms like AWS assume individual machines are disposable. Instead of building one ultra-reliable server, you run your service across dozens of commodity machines and design the software to handle any of them disappearing at any moment. This approach is cheaper, more flexible, and scales better. When a VM dies, the orchestrator spins up a replacement in seconds. The application never notices.
The practical implication: do not design your system around the assumption that any particular machine will be alive tomorrow. Design it so that any machine can disappear and the system keeps serving users. This mindset shift from "protect the machine" to "protect the service" is the foundation of cloud-native architecture.
An important nuance: hardware faults are becoming more common as systems grow, not because hardware is getting worse, but because probability compounds with scale. If each server has a 0.1% chance of failing in a given day, a cluster of 10 servers has roughly a 1% chance of experiencing at least one failure. At 1,000 servers, hardware failure is a daily occurrence. At 10,000 servers, you get multiple failures per day. The math is simple, but the design implications are profound: your system must treat hardware failure as a routine event, not an emergency.
Software Faults
Software faults are more dangerous than hardware faults because they tend to be correlated. A hardware failure is usually independent: one disk dying does not cause another disk to die. But a software bug can hit every node in a cluster simultaneously. A leap-second bug once crashed every Linux kernel of a certain version across the globe at the same instant. A cascading failure in one service can trigger timeouts in every service that depends on it, turning a single slow database query into a system-wide outage.
Software faults lurk for months or years, waiting for an unusual input or a specific set of conditions to trigger them. A memory leak that takes weeks to exhaust available RAM. A race condition that only manifests under high load on the third Tuesday of the month. An assumption that a downstream service always responds in under 100ms that breaks when that service deploys a new version. A null pointer that only triggers when a user's name contains a Unicode character the developer never tested.
The defenses against software faults are fundamentally different from hardware. You cannot just add redundancy because the same bug exists on every replica. Instead you need: thorough testing (including chaos engineering that deliberately injects failures), process isolation (a crash in one process should not take down the host), monitoring with alerting (detect slow degradation before it becomes an outage), and careful design of failure domains so that a bug in one subsystem cannot cascade to others.
One powerful technique is the bulkhead pattern, borrowed from ship design. A ship's hull is divided into watertight compartments so that a breach in one compartment does not sink the entire ship. Similarly, you can isolate services into failure domains: a bug in the recommendation engine should not affect the checkout flow. Thread pools, separate process groups, and independent deployment units all serve as bulkheads.
Another defense is the timeout-and-retry pattern with exponential backoff. When a service calls a dependency that is responding slowly, it should not wait indefinitely. Set a timeout (say, 500ms), and if the dependency has not responded, fail the call and either retry with a longer wait or return a degraded response to the user. Without timeouts, a single slow dependency ties up threads in the calling service, which then cannot serve other requests, which causes its callers to time out, cascading the failure up the entire call chain. Circuit breakers formalize this: after a threshold of failures, the circuit "opens" and all subsequent calls fail immediately without even attempting the downstream call, giving the failing service time to recover.
That is the one-sentence version, and it is all this lesson needs. The three states, the thresholds that move a breaker between them, and the interaction with the pool of connections the breaker is protecting are developed in connection pooling and resource management; the breaker's place alongside timeouts, retry budgets and bulkheads is in designing for failure. The reason it appears here at all is that a circuit breaker is the clearest case of reliability and maintainability pulling against each other: it makes the system survive a dependency outage, and it adds a second state machine that can itself be misconfigured, as the closing section of this lesson argues.
The combination of these patterns creates a layered defense:
Graceful degradation deserves special attention. When a non-critical dependency fails (the recommendation engine is down), the system should not return an error to the user. Instead, show a default recommendation list, or hide the recommendation section entirely. The user gets a slightly worse experience, but the core functionality (browsing, searching, purchasing) remains available. Netflix calls this "the fallback pattern": every service call has a predefined fallback response for when the service is unavailable.
The critical difference between hardware and software faults is correlation. Hardware failures are mostly independent, so redundancy works. Software failures are often correlated, hitting every replica simultaneously, so redundancy alone cannot save you. This is why Netflix invented Chaos Monkey: you need to test that your fault tolerance actually works by injecting real faults in production.
Human Errors
Humans are the least reliable component in any system. Studies of large internet services found that configuration errors by operators were the leading cause of outages, far exceeding hardware or software faults. This is not because operators are careless. It is because humans operate in complex environments under time pressure, and the interfaces they work with are often error-prone by design.
Consider the difference between a configuration file where one typo takes down production and an infrastructure-as-code pipeline where every change goes through code review, automated validation, and staged rollout. The human making the change is the same person. The error rate is different because the system design is different.
Designing for human error means building systems that minimize the opportunity for mistakes and limit the blast radius when mistakes happen. The distinction matters: minimizing opportunity (making mistakes harder to make) is prevention; limiting blast radius (reducing damage when mistakes happen) is containment. You need both, because prevention alone is never 100% effective.
Specific strategies include:
- Sandbox environments where operators can experiment with production-like conditions without risk to real users
- Rollback mechanisms that undo a bad deploy in seconds rather than hours, ideally automated when health checks fail
- Detailed monitoring that surfaces unexpected behavior early, before a small anomaly becomes a cascading failure
- Gradual rollout practices (canary deploys, feature flags) that expose changes to a small percentage of traffic before going wide
- Immutable infrastructure where servers are replaced rather than modified in place, eliminating "snowflake server" drift
The best organizations treat human error not as a character flaw to punish but as a design problem to solve. If an operator can make a single mistake that takes down production, the system's design is at fault for allowing it, not the operator. Post-incident reviews should focus on "how do we make this class of mistake impossible?" not "who made the mistake?"
A concrete example: Google's approach to safe configuration changes requires that any config change goes through a pipeline that validates the syntax, applies the change to a small canary subset, monitors key metrics for anomalies, and only then rolls the change out progressively to the full fleet. An engineer cannot push a config change directly to all servers even if they wanted to. The system makes the dangerous path harder than the safe path. This is designing for human error at its best.
Measuring Reliability: SLOs and Error Budgets
Reliability is not binary. No system is 100% reliable, and attempting perfect reliability is both impossible and counterproductive (you would never ship new features for fear of introducing bugs). Instead, modern teams define Service Level Objectives (SLOs) that quantify how reliable the system needs to be.
An SLO might state: "99.9% of requests will return a successful response within 500ms." This means you are explicitly accepting that 0.1% of requests can fail or be slow. Over a month, that is about 43 minutes of allowed downtime, your error budget.
The error budget concept transforms reliability from a vague aspiration into a concrete engineering tradeoff. When the error budget is healthy (few incidents), the team ships features faster, accepts more risk, and deploys more aggressively. When the error budget is nearly exhausted (too many incidents this month), the team slows down, focuses on stability work, and defers risky changes. This makes the reliability-vs-velocity tradeoff explicit and data-driven rather than a recurring argument between engineering and product teams.
The key insight is that reliability is a feature with a cost, and like any feature, it should be budgeted based on user needs, not maximized without limit. A 99.99% SLO costs 10x more engineering effort than a 99.9% SLO, and for many products the users cannot even perceive the difference.
The budget itself is one multiplication. Over a window of length , an objective of leaves
so 99.9% over thirty days is minutes, and 99.99% is 4.32. The number that actually drives operational decisions is the second derivative of that: the burn rate,
which is dimensionless and says how many times faster than sustainable you are spending. exhausts the budget exactly at the end of the window. burns 2% of a thirty-day budget in one hour, which is the threshold Google's SRE workbook uses for a page, paired with a slower six-hour window at to catch the grinding degradations that never trip the fast alarm. Alerting on burn rate rather than on raw error rate is what stops a 5% error spike lasting forty seconds from waking someone at 3 a.m.
To put the numbers in perspective:
Each additional "nine" roughly 10x the engineering investment. Most applications do not need more than three nines. A batch processing pipeline might only need two. A payment processing system might need four. The SLO should match the user's tolerance for failure, not the engineer's ambition for perfection.
A common mistake is setting SLOs without considering dependencies. If your system depends on a cloud database with a 99.95% SLO and a third-party payment API with a 99.9% SLO, your system's maximum achievable availability is the product of its dependencies' availabilities: roughly 99.85%. Setting your own SLO at 99.99% is meaningless if your dependencies cannot support it. Realistic SLOs account for the entire dependency chain.
In practice, SLOs are most effective when they are visible to both engineering and product teams. Display error budget consumption on a shared dashboard. When the error budget is healthy, product managers can push for faster feature delivery. When it is nearly exhausted, engineering can slow down and focus on stability without needing to justify the decision politically. The error budget becomes a shared language between teams that have historically different incentives.
SLOs should be measured from the user's perspective, not the server's perspective. A server might report 100% uptime while users experience errors because the load balancer is misconfigured, the CDN is serving stale 404 pages, or DNS is resolving to the wrong IP. User-facing SLOs measure success at the edge: did the user's request return a correct response within the latency target? This often requires synthetic monitoring (periodic test requests from multiple geographic locations) in addition to server-side metrics.
One final note on reliability measurement: distinguish between planned and unplanned downtime. A 30-minute maintenance window at 3 AM (communicated in advance, during low traffic) is very different from a 30-minute outage during peak hours caused by a failed deployment. Both consume error budget, but their impact on user trust is vastly different. Some teams track two separate error budgets: one for total availability and one for unplanned availability.
Scalability is not a binary label you slap on a system. You do not say "this system is scalable" the way you say "this system uses PostgreSQL." Scalability is a question: "If load increases in a specific way, what are our options for coping with the growth?" The answer depends entirely on what kind of load you have, how you measure performance, and what tradeoffs you are willing to make.
A common trap in system design interviews and real architecture discussions alike is jumping to scaling solutions without first defining the scaling problem. "We need to handle more traffic" is not a problem statement. "Our p99 latency exceeds 2 seconds when we pass 5,000 RPS because the database connection pool saturates at 4,800 concurrent queries" is a problem statement. The first leads to generic solutions (add more servers). The second leads to targeted solutions (increase connection pool size, add read replicas for read-heavy queries, implement connection pooling with PgBouncer).
Describing Load: Load Parameters
Before you can talk about scaling, you need a precise vocabulary for describing load. Vague statements like "we get a lot of traffic" are useless for capacity planning. Load parameters are the specific numbers that define your system's workload. Without them, scaling conversations devolve into guesswork.
Common load parameters include:
The key insight is that every system has a small number of load parameters that dominate its scaling challenges. Twitter's defining parameter was fan-out: when a celebrity with 30 million followers tweets, that single write must become 30 million reads in home timelines. The solution was to precompute timelines for most users (fan-out on write) but fall back to fan-out on read for celebrities. You cannot design this solution without first identifying fan-out as the critical load parameter.
Another way to think about load parameters is as the dimensions along which your system can break. A payment processing system might handle 1,000 RPS easily but choke at 100 RPS if each request involves a 3-second external API call to a bank. The RPS number alone says nothing. The combination of RPS and per-request external latency tells you everything. Similarly, a chat application that handles 10,000 messages per second might be fine until a single group chat has 50,000 members, because fan-out per message is the parameter that matters, not total message throughput.
The Load Curve and the Knee
Every system has a characteristic load curve: as load increases, performance stays relatively stable until it hits a threshold, then degrades sharply. This inflection point is called the "knee" of the load curve. Below the knee, the system absorbs additional load gracefully. Above the knee, small load increases cause disproportionate performance degradation.
Why does the knee exist? Because systems have finite resources (CPU, memory, connections, disk IOPS), and as utilization approaches capacity, contention increases nonlinearly. At 50% CPU utilization, processes rarely wait for a core. At 90% utilization, processes frequently queue, and each queued process adds latency that cascades to everything waiting on it. The same effect happens with database connections, thread pools, and network buffers.
The practical implication is that you should never run a system at more than 70-80% of its maximum capacity in production. The remaining 20-30% is your headroom for traffic spikes, garbage collection pauses, and the occasional runaway query. If your system routinely operates above 80% of any resource, you are living on the knife's edge of the load curve, and any spike will push you past the knee into degraded performance.
Finding the knee requires load testing. Tools like Locust, k6, and Gatling let you simulate increasing load against your system and observe how latency, throughput, and error rate respond. The test protocol is simple: start at 10% of expected load, increase by 10% every 5 minutes, and record metrics at each step. Plot latency versus load and the knee becomes visible: the point where the latency curve bends upward. Your production capacity target should be 70-80% of the load at the knee, giving you margin for unexpected spikes.
Describing Performance: Throughput, Latency, and Percentiles
Once you know your load, you need to measure how the system performs under that load. Two systems can handle 10,000 RPS but feel completely different to users if one responds in 5ms and the other in 500ms.
Throughput is the number of operations processed per unit time. A batch processing system cares about throughput: how many records can it process per hour? A MapReduce job processing 100 TB of logs cares about total time to completion, not how fast any individual record is processed. Throughput and latency are often inversely correlated: batching requests together increases throughput (fewer round-trips, better amortization of fixed costs) but increases latency for individual requests (each request waits for the batch to fill).
Latency is the duration a single operation takes from the client's perspective. An online system cares about latency: how long does the user wait for a response? Note the distinction between latency (the full round-trip time the user experiences, including network, queuing, and processing) and service time (just the processing time on the server). A request might have 2ms of service time but 200ms of latency because it waited 198ms in a queue.
Never use averages to describe latency. An average of 100ms might mean that 99% of requests complete in 50ms and 1% take 5 seconds. The average hides the pain. Always use percentiles: p50 (median), p99 (worst 1%), and p999 (worst 0.1%). The users who experience tail latency are often your most valuable: they have the most data, the most activity, and the most complex queries.
Percentiles expose the distribution of latency. The p50 (median) tells you the typical experience. The p99 tells you the worst-case experience for 1 in 100 requests. The p999 tells you the experience that 1 in 1,000 users hits. Amazon has found that a 100ms increase in p99.9 latency correlates with a 1% decrease in revenue, because the slowest requests belong to the customers with the most items in their cart.
Why does tail latency matter so much? Because in a microservice architecture, a single user request often fans out to multiple backend services. If each of 5 backend calls has a 1% chance of being slow (p99 = 2 seconds), the probability of at least one being slow is roughly 5%. That means 5% of user requests hit the 2-second tail latency, even though each individual service is "fast" 99% of the time. This is called tail latency amplification, and it gets worse as you add more services to the call chain.
If the calls are independent and each is slow with probability , the chance that a request touching services hits at least one slow call is
At , five calls give 4.9% and twenty calls give 18.2%. Nearly one user request in five lands on a two-second tail while every service dashboard shows a healthy 99th percentile. The exponent is the part worth internalizing: the fan-out, not the individual service, is the variable you are actually tuning, and a service mesh that turns one call into three retries has tripled without anyone editing an SLO.
The practical consequence: if your request path touches 20 services and each has a p99 of 100ms, the effective p99 for the entire request is far worse than 100ms. Reducing individual service p99 from 100ms to 10ms might matter more than reducing p50 from 5ms to 1ms, because the tail compounds across services.
How do you actually measure percentiles in production? You cannot store every request latency forever. Instead, systems use approximation algorithms like t-digest or HDR Histogram that maintain a compact, streaming summary of the distribution. These data structures let you query any percentile with bounded error using a fraction of the memory that storing every value would require. Most observability platforms (Datadog, Prometheus with histograms, AWS CloudWatch) compute percentiles this way under the hood.
A common mistake is computing percentiles by averaging percentiles from individual servers. If you have 10 servers each reporting their p99, averaging those 10 numbers does not give you the global p99. Percentiles are not additive. To get the global p99, you need to aggregate the raw histograms from all servers and compute the percentile from the combined distribution. This is why centralized metrics collection matters: you need the full picture, not summaries of summaries.
It is worth seeing how far off that shortcut goes, because the intuition is that averaging ten p99s should be roughly right:
One degraded server out of ten, and the averaged figure understates the real tail by a factor of 2.1. The direction is not accidental. Averaging pulls the outlier toward nine healthy servers, but the global p99 is drawn from the slowest 1% of all one million requests, and the degraded server contributes a disproportionate share of them. The more skewed the fleet, the worse the shortcut gets, which means it fails hardest exactly during the incident you built the dashboard for. The fix is to export histogram buckets rather than computed percentiles, so the aggregation happens over the distribution: Prometheus histogram_quantile over summed buckets is correct, avg(p99) over a recording rule is not.
Another subtlety: head-of-line blocking can make latency measurements misleading. If a single slow request blocks a thread or connection, all subsequent requests queued behind it appear slow too, even though their actual processing time is fast. The measured p99 reflects queuing time, not service time. Identifying whether tail latency comes from slow processing or from queuing contention requires tracing individual requests through the system, not just measuring end-to-end latency at the load balancer.
A final performance concept worth internalizing is coordinated omission: when your load testing tool waits for the previous request to complete before sending the next one, it under-reports latency because it stops generating load during slow periods. A real user does not stop clicking because the server is slow. Real load generators should send requests at a fixed rate regardless of response time. Tools like wrk2 and Gatling handle this correctly; simpler tools often do not. Ignoring coordinated omission can make your p99 measurements look 10x better than reality.
Vertical vs. Horizontal Scaling
When load exceeds capacity, you have two fundamental options.
Vertical scaling (scaling up) means buying a bigger machine: more CPU cores, more RAM, faster disks. The advantage is simplicity. Your application code does not change. One machine, one process, no distributed coordination. There are no network partitions between components because everything communicates through shared memory. The disadvantages are cost (hardware pricing is superlinear: doubling RAM costs more than double the price) and hard limits (the biggest machine money can buy still has a ceiling, and you cannot buy a server with 100 TB of RAM).
Horizontal scaling (scaling out) means adding more machines and distributing the load across them. The advantage is that there is no practical ceiling: you can always add another machine to the fleet. The disadvantages are significant complexity. Your data must be partitioned across machines. Your requests must be routed to the right machine. Your state must be coordinated across replicas. Distributed systems introduce failure modes (network partitions, clock skew, split-brain) that do not exist on a single machine. You now need consensus protocols, distributed transactions, or eventual consistency, each with its own tradeoffs.
In practice, most systems use a hybrid approach. Keep things on a single machine as long as possible (vertical scaling is simpler and cheaper up to a point), then go horizontal when you hit the cost-performance crossover or the single-machine ceiling. The crossover point is system-specific: a CPU-bound service hits it sooner than a memory-bound service because CPU capacity maxes out faster than memory capacity on modern hardware.
A useful mental model: vertical scaling buys you time, horizontal scaling buys you capacity. When your database is running hot, upgrading to a larger instance class takes an hour and buys you 6 months of headroom. That 6 months is time you can use to architect a proper sharding strategy. Teams that jump straight to horizontal scaling without first buying time through vertical scaling often end up with rushed, poorly designed distributed systems that create more problems than they solve.
The cost model also matters. Cloud providers price compute superlinearly at the top end: the largest instance types cost significantly more per unit of CPU/RAM than mid-tier instances. At some point, running 4 medium instances is cheaper than running 1 extra-large instance with the same total capacity. This economic crossover point varies by cloud provider and instance family, but it typically occurs around the top 2-3 instance tiers. For cost-sensitive workloads, this crossover point, not a technical limit, is what drives the move to horizontal scaling.
There is also a resilience argument for horizontal scaling: 4 medium instances can survive the loss of one and continue serving traffic at 75% capacity, while 1 extra-large instance is a single point of failure. The hybrid strategy of vertical scaling for simplicity with horizontal scaling for resilience gives you both: a small number of medium-sized machines, each large enough to be efficient but numerous enough that losing one is not catastrophic.
There is an important distinction between stateless and stateful services when scaling horizontally. Stateless services (API servers, web frontends, computation workers) scale horizontally with minimal complexity: just add more instances behind a load balancer. Each instance handles any request independently. Stateful services (databases, caches, message brokers) are much harder to scale horizontally because you must partition the data, maintain consistency guarantees, and handle the case where a client's data lives on a different machine than the one receiving the request.
A common scaling anti-pattern is premature horizontal scaling of stateful services. A team adds database sharding at 1,000 users because "we want to be ready for growth." Now every query that joins data across shards requires a scatter-gather pattern. Simple features become 10x harder to implement. The team spends more time fighting the distributed database than building product features. Meanwhile, a single PostgreSQL instance with proper indexing and a connection pooler could have handled 100,000 users. The rule of thumb: exhaust vertical scaling and query optimization for your database before introducing horizontal partitioning.
The Scaling Sequence
In practice, most systems follow a predictable scaling sequence as they grow. Understanding this sequence helps you anticipate what comes next and prepare for it without over-investing prematurely.
Phase 1 (0-10K users): Single server. One application server, one database server. Optimize queries, add indexes, enable connection pooling. This carries you further than most people expect.
Phase 2 (10K-100K users): Separate concerns. Move the database to a dedicated server. Add a CDN for static assets. Add a reverse proxy (Nginx) for SSL termination and static file serving. Introduce caching (Redis or Memcached) for frequently accessed data.
Phase 3 (100K-1M users): Horizontal application tier. Run multiple application server instances behind a load balancer. Add database read replicas for read-heavy queries. Implement session storage in Redis so any app server can handle any request.
Phase 4 (1M-10M users): Specialize and partition. Introduce message queues for async processing. Separate read and write database paths. Begin database sharding for the largest tables. Add monitoring, alerting, and automated scaling.
Phase 5 (10M+ users): Distributed everything. Multi-region deployment. Service-oriented architecture or microservices. Custom infrastructure for specific bottlenecks. Dedicated teams for infrastructure, performance, and reliability.
Each phase builds on the previous one. The mistake is jumping to Phase 5 at Phase 1 traffic levels. Every phase has a cost in complexity, operational overhead, and engineering time. Pay that cost only when the current phase's limits are measurably reached.
Notice that caching appears early in this sequence (Phase 2), before horizontal scaling. This is intentional. A well-implemented cache (90%+ hit rate on read-heavy endpoints) can reduce database load by 10x, which is equivalent to 10x horizontal scaling of the database tier but with far less complexity. For read-heavy workloads, caching is the single highest-leverage scaling technique available. For write-heavy workloads, caching helps less, and you need to move to Phase 4 (sharding) sooner.
Elastic Scaling and Auto-Scaling
A modern refinement of horizontal scaling is elastic scaling: automatically adding or removing instances based on current load. Cloud platforms make this straightforward for stateless services. Set a CPU utilization target (say, 60%), define minimum and maximum instance counts, and let the auto-scaler adjust capacity in real time.
The benefit is cost efficiency: instead of provisioning for peak load 24/7, you provision for average load and let the auto-scaler handle spikes. A service that peaks at 10x its baseline during a daily batch job can run on 3 instances most of the day and scale to 30 instances for the 2-hour peak window. Without auto-scaling, you would pay for 30 instances around the clock.
The challenge with elastic scaling is the cold-start penalty. New instances take time to initialize: downloading container images, establishing database connections, warming caches, and loading models into memory. If your service takes 3 minutes to start and a traffic spike arrives suddenly, the auto-scaler cannot respond fast enough. Solutions include pre-warming (keeping a buffer of idle instances), predictive scaling (scaling up before anticipated peaks based on historical patterns), and reducing startup time through smaller containers and lazy initialization.
Elastic scaling also requires careful metric selection. Auto-scaling on CPU utilization works well for compute-bound workloads but is misleading for I/O-bound services where CPU sits at 20% while the service is saturated waiting on database responses. For I/O-bound services, scale on request queue depth or response latency instead. For memory-bound services (caches, in-memory databases), scale on memory utilization. The wrong scaling metric leads to either under-provisioning (scaling triggers never fire despite degraded performance) or over-provisioning (scaling triggers fire too aggressively, wasting money on instances that sit idle).
One more consideration: auto-scaling down is as important as scaling up, but riskier. Scaling up adds capacity, which is always safe. Scaling down removes capacity, which means existing connections must be drained, in-flight requests must complete, and the remaining instances must absorb the redistributed load without exceeding their own capacity. A common practice is to scale down slowly (one instance at a time, with a cooldown period between removals) even though scaling up happens quickly (add multiple instances simultaneously).
A related concept is scheduled scaling: pre-provisioning additional capacity for known traffic patterns. An e-commerce platform that sees 3x traffic on Black Friday should not rely on reactive auto-scaling alone. Instead, schedule additional instances to be running before the traffic spike arrives, so there is no cold-start delay when the spike hits. Combine scheduled scaling (for predictable patterns) with reactive auto-scaling (for unexpected spikes) for the most robust approach.
The overarching principle of scaling is to match capacity to demand with minimal waste and minimal risk. Too little capacity causes outages and poor user experience. Too much capacity wastes money. The tools described in this section (load parameters, percentiles, vertical/horizontal scaling, elastic scaling) are all in service of finding the right capacity at the right time.
In interviews, when asked about scaling, always start by identifying the bottleneck. Is it CPU, memory, disk I/O, or network? A CPU-bound service scales horizontally by adding stateless workers behind a load balancer. A memory-bound service (like a large cache) scales by partitioning data across machines. A disk I/O-bound service scales by moving to SSDs first (vertical) before sharding (horizontal). The scaling strategy depends on the bottleneck, not a generic recipe.
The majority of the cost of software is not in its initial development but in its ongoing maintenance: fixing bugs, keeping systems running, investigating failures, adapting to new platforms, repaying technical debt, and adding new features. Studies consistently show that 60-80% of total software cost occurs after the initial release. Yet when engineers design systems, they spend 90% of their attention on the initial build and 10% on what happens for the next five years. Maintainability flips this ratio in your thinking.
Maintainability is not one property. It is three distinct qualities that serve different stakeholders and require different design decisions. Operations teams care about operability. New engineers care about simplicity. Product teams care about evolvability. A well-designed system serves all three.
Operability: Making Life Easy for Operations
A system with good operability is one that operations teams can run without heroics. This means the system provides visibility into its internal state (metrics, logs, distributed traces), supports automation (self-healing, auto-scaling), avoids hard-coded dependencies on specific machines or IP addresses, gracefully handles machine failures without manual intervention, and has comprehensive documentation for operational runbooks.
The difference between good and bad operability is the difference between "the on-call engineer gets paged, runs a documented playbook, and resolves the issue in 10 minutes" and "the on-call engineer gets paged, spends 3 hours reading code to understand what the system is even doing, and fixes it by restarting everything and hoping for the best." The first scenario is sustainable. The second leads to burnout, attrition, and increasingly fragile systems as institutional knowledge leaves with departing engineers.
Good operability is not glamorous work, but it is what separates systems that run smoothly for years from systems that require constant babysitting. Invest in observability early in the project lifecycle, not after the first outage. Specifically, this means structured logging (JSON logs with request IDs, not printf statements), metrics dashboards that show both system health and business metrics, alerting thresholds based on symptoms (error rate, latency) not causes (CPU utilization), and distributed tracing that follows a request across service boundaries.
A practical test for operability: can you answer "why is the system slow right now?" within 5 minutes using only your dashboards and logs, without reading any application code? If yes, your operability is solid. If you need to SSH into servers and grep log files, you have work to do.
Another dimension of operability is automation. Systems that require manual steps for common operations (deploys, scaling, failover, database migrations) accumulate risk with every manual step. Each manual step is an opportunity for human error. Automate the routine: rolling deploys, automated health checks, auto-scaling policies based on CPU or queue depth, and automated database backups with regular restore tests. The goal is to reduce the set of manual operations to only truly novel situations that require human judgment.
Capacity planning is an often-overlooked aspect of operability. A well-operated system has projections for when current capacity will be exhausted, based on growth trends. "At current growth rate, our database will hit disk limits in 6 months" is the kind of forward-looking statement that operational maturity enables. Without capacity planning, teams discover limits in production when the system degrades, turning a predictable problem into an emergency.
Finally, documentation is an operability feature, not a nicety. Runbooks that describe how to diagnose and resolve common issues. Architecture decision records (ADRs) that explain why the system is designed the way it is. On-call handoff documents that brief the next rotation on ongoing issues. These artifacts reduce the knowledge gap between the engineer who built the system and the engineer who must operate it at 3 AM during an outage.
The distinction between good and great operability is proactive versus reactive. Good operability detects and responds to problems quickly. Great operability prevents problems by surfacing trends, running pre-mortem exercises ("what would happen if this component failed?"), and conducting regular game days where the team practices incident response against simulated failures. The teams with the best operability invest 20-30% of engineering time in operational improvements, not as a one-time project but as a continuous practice.
Simplicity: Managing Complexity
Complexity is the enemy of maintainability. Not the inherent complexity of the problem domain (a financial trading system is inherently complex because financial regulations are complex) but the accidental complexity introduced by poor design decisions. Accidental complexity includes tangled dependencies where changing one module requires understanding ten others, inconsistent naming conventions across the codebase, workarounds for workarounds that nobody remembers the reason for, and special cases that exist because nobody refactored when requirements changed.
The best tool for managing complexity is abstraction. A well-designed abstraction hides implementation details behind a clean interface, letting you reason about each layer independently. TCP hides packet retransmission from application code. SQL hides storage layout from query logic. Good abstractions do not just make code shorter. They make it possible for a new engineer to contribute to the system without understanding every layer beneath the one they are working on.
The symptom of accidental complexity is fear. When engineers are afraid to change a module because they do not understand its side effects, that is accidental complexity. When a "simple" feature takes 3 months because it touches 15 services, that is accidental complexity. When onboarding a new engineer takes 6 months instead of 2 weeks, that is accidental complexity. When the team decides to rewrite the system from scratch rather than evolve it, that is the terminal stage of accidental complexity.
How do you fight accidental complexity? Ruthless refactoring when you see it, not in a big bang rewrite but incrementally. Boy Scout rule: leave the code cleaner than you found it. When you touch a module, fix its naming, extract its dependencies, add its tests. Over time, the codebase improves section by section. The key is that this incremental improvement must be culturally valued and time-budgeted, not treated as an afterthought squeezed in between feature work.
A useful metric for tracking complexity is coupling. How many other modules does a given module depend on? How many modules depend on it? A module with 20 incoming and 15 outgoing dependencies is a complexity hotspot: changing it risks breaking everything, and everything that changes might break it. Visualizing dependency graphs and tracking coupling metrics over time reveals whether your refactoring efforts are actually reducing complexity or just moving it around.
Another useful lens is cognitive load: how much does an engineer need to hold in their head to make a change? A function with 5 parameters, 3 nested conditionals, and 2 global side effects requires the engineer to track many things simultaneously. Breaking it into smaller functions with explicit inputs and outputs reduces cognitive load, even if the total line count increases. Code is read 10x more often than it is written, so optimizing for readability beats optimizing for brevity every time.
There is a paradox in simplicity: sometimes making the system simpler for one audience makes it more complex for another. A well-designed abstraction simplifies life for the engineers who use it but adds complexity for the engineers who maintain it. A message queue simplifies application code (just publish and subscribe) but adds an infrastructure component that must be monitored, scaled, and recovered from failures. The net effect is still positive because the abstraction encapsulates complexity in one place (the queue implementation) rather than spreading it across every service that needs async communication. But it is a reminder that simplicity is about managing complexity, not eliminating it.
One practical technique for managing complexity is the bounded context from domain-driven design. A bounded context defines a clear boundary within which a particular domain model is consistent and self-contained. Inside the boundary, terms have precise meanings and the model is internally coherent. At the boundary, you translate between different models. This prevents the "big ball of mud" problem where every concept leaks across every module, and a change to the definition of "user" in one context breaks 20 other contexts that use the same shared model with slightly different assumptions.
The number one complexity killer is shared mutable state. When multiple components read and write the same data with unclear ownership, every change is a potential conflict. The fix is clear data ownership: for each piece of data, exactly one component is the authoritative writer, and all other components either read from it or receive events about changes. This pattern (sometimes called "single writer principle") dramatically reduces the complexity of reasoning about data consistency.
Technical debt is a useful metaphor but often misused. Real technical debt is a deliberate tradeoff: "We are taking a shortcut here because time-to-market matters more right now, and we will pay it back next quarter." That is a rational decision. What most teams call "tech debt" is actually reckless negligence: shortcuts taken without awareness of the cost, never planned to be repaid, accumulating silently until the system becomes unmaintainable. The first kind is a strategic tool. The second is an organizational failure. Distinguishing between them matters because the first can be managed through prioritization, while the second requires cultural change.
Evolvability: Making Change Easy
Requirements change. Markets shift. Regulations get introduced. Users request features nobody anticipated. Technology evolves and what was a reasonable choice 3 years ago becomes a liability. A maintainable system is one that accommodates change without requiring a ground-up rewrite.
Evolvability is closely linked to simplicity: simpler systems are easier to modify because each change affects fewer components. But evolvability also requires intentional design choices that go beyond "clean code." Clean API boundaries between components mean you can replace the internals of one component without affecting its consumers. Modular architecture means new features can be added as new modules rather than woven into existing code. Comprehensive automated tests mean refactoring is safe because regressions are caught immediately, not by users in production.
Consider the difference between adding a new payment method to a system with a well-defined PaymentProvider interface versus one where Stripe API calls are scattered across 40 files. The first system requires implementing a new class that satisfies the interface. The second requires finding and modifying every file that touches payment logic, with no guarantee you found them all.
Version control discipline also supports evolvability. Small, focused commits with clear messages make it possible to understand why a change was made months later. Feature branches with code review catch design problems before they merge. Continuous integration that runs the full test suite on every commit catches regressions immediately. These practices are not just "process overhead." They are the mechanisms that make change safe, and safe change is the prerequisite for rapid change.
A subtle but important design choice for evolvability is the dependency inversion principle: high-level modules should not depend on low-level modules. Both should depend on abstractions. In practice, this means your business logic should not import your database driver directly. Instead, it should depend on a repository interface, and the database driver implements that interface. When you need to switch from PostgreSQL to DynamoDB, or add a caching layer, or run tests against an in-memory store, you implement a new version of the interface without touching the business logic. The upfront cost of this abstraction is small. The long-term evolvability benefit is enormous.
Another design choice that enables evolvability is backward-compatible API design. When internal or external APIs evolve, new versions should support old clients without breaking them. Techniques include adding new fields to responses without removing old ones, using API versioning (v1, v2) for breaking changes, and defaulting new required parameters to sensible values for old clients. This allows different components of the system to evolve at different rates, which is essential when you have many teams or many external consumers who cannot all upgrade simultaneously.
The test for evolvability is practical: can a new engineer, reading the codebase for the first time, make a meaningful change within their first week? If yes, the system is evolvable. If they need months of context before touching anything, the system has an evolvability problem regardless of how well-designed the original architecture was.
The Cost of Neglecting Maintainability
The consequences of poor maintainability are not immediate. They compound over months and years. In the first year, a messy codebase feels fast because there are no abstractions to slow you down. By year two, feature velocity has dropped by half because every change requires understanding tangled dependencies. By year three, the team is spending 70% of its time on bug fixes and operational firefighting, with only 30% available for new features. By year five, the team proposes a full rewrite because the codebase is "unmaintainable."
The full rewrite is almost always a mistake. It takes longer than estimated (typically 2-3x), it must replicate years of bug fixes and edge cases that are embedded in the old code but not documented, and the business must maintain two systems in parallel during the migration. The better path is incremental improvement: the strangler fig pattern for monolith-to-microservice migrations, systematic addition of tests to legacy code, and gradual extraction of well-defined modules behind clean interfaces.
The full rewrite trap is also driven by a cognitive bias: engineers overestimate how well they understand the old system and underestimate how much accidental complexity they will reintroduce in the new one. The new system starts clean, but under the same time pressure and organizational incentives that created the original mess, it follows the same trajectory. Three years later, someone proposes rewriting the rewrite.
Joel Spolsky called this "the single worst strategic mistake that any software company can make": throwing away working code and starting from scratch. The working code, ugly as it may be, embodies years of bug fixes, edge case handling, and hard-won knowledge about real-world conditions that no specification document captures. The strangler fig pattern, where you incrementally replace components of the old system with new implementations while both run in parallel, preserves this institutional knowledge while still achieving modernization.
The strangler fig approach works as follows: identify a well-defined subsystem at the boundary of the monolith (an API endpoint, a background job, a data pipeline). Build the replacement as an independent service. Route traffic to the new service for that subsystem while the old code remains available as a fallback. Once the new service is stable, remove the old code. Repeat for the next subsystem. This is slower than a big-bang rewrite but dramatically safer, and each step delivers measurable value.
The takeaway: maintainability is not a luxury for teams that have time. It is an investment that determines whether your team has time in the future. Every shortcut taken today is a tax paid on every change tomorrow.
Reliability, scalability, and maintainability are not independent goals you can maximize in isolation. They interact, and optimizing for one often creates tension with another. The art of system design is finding the right balance for your specific context, your team size, your user base, and your business stage, not maximizing all three simultaneously.
Reliability vs. Scalability
Scaling a system often introduces new reliability risks. A single-server application has a simple failure model: the server is up or down. There are no partial failures, no network partitions between components, and no consensus problems. Scale it to 50 servers across 3 availability zones, and now you have network partitions where servers cannot communicate, partial failures where 3 of 50 servers are down but the rest are fine, clock skew where servers disagree on the current time, and the need for distributed consensus to coordinate state changes. Every horizontal scaling step adds failure modes that did not exist before.
Conversely, reliability requirements can constrain scaling options. If your system requires strong consistency (every read reflects the most recent write), you cannot simply add read replicas that serve stale data. You need synchronous replication or consensus protocols like Raft or Paxos, both of which limit throughput and increase latency compared to eventual consistency. The CAP theorem is not just a theoretical concept. It is a real engineering constraint that forces you to choose between consistency and availability during network partitions. Your scaling strategy must respect your consistency requirements.
A real-world example: an e-commerce platform needs its inventory count to be consistent (you cannot sell an item that is already sold). This consistency requirement means you cannot freely replicate the inventory database across regions with eventual consistency, because two regions might both sell the last item. The solution is typically to partition inventory by geographic region or warehouse, with strong consistency within each partition. This limits scaling flexibility but preserves the correctness guarantee that the business depends on. The engineering tradeoff is explicit: we sacrifice some scaling options to maintain the reliability property (correctness) that matters most.
The inverse is also common: teams that over-specify consistency requirements. A social media feed does not need strong consistency. If a user sees a post 500ms later than another user, nobody notices or cares. But if you implement the feed with synchronous replication for "reliability," you have artificially limited your scaling options and added latency for no user-visible benefit. Match your consistency requirements to what users actually need, not to what feels theoretically correct.
A useful exercise when evaluating reliability-scalability tradeoffs: for each data store or service in your system, ask "what is the worst thing that happens if this is unavailable for 5 minutes? What about 5 seconds of stale data?" For a user profile service, 5 seconds of stale data is invisible to users. For an account balance service, even 1 second of stale data could result in an overdraft. The answers determine whether you need synchronous replication (strong consistency, lower availability) or asynchronous replication (eventual consistency, higher availability and throughput).
Many systems contain a mix: the user profile service uses eventual consistency (high availability, lower latency) while the account balance service uses strong consistency (correctness guarantee, higher latency). This heterogeneous approach is not a compromise. It is the optimal design, because different data has different consistency requirements, and treating all data the same wastes either correctness or performance.
Scalability vs. Maintainability
The fastest way to scale is often the least maintainable. Highly optimized code that squeezes every microsecond out of a hot path is harder to understand, harder to modify, and harder for new engineers to contribute to. A hand-tuned caching layer with custom eviction logic, shard-aware routing, and async replication handles load beautifully but requires deep system knowledge to operate and modify. When the original author leaves the company, that component becomes a black box that everyone is afraid to touch. This phenomenon has a name: bus factor. If only one person understands a critical system component, the project is one resignation away from a crisis.
On the other hand, a clean and simple codebase that any engineer can understand in a week might not handle your scale. The art is knowing when to introduce complexity (and accepting the maintenance cost) versus when to pay for simpler design with more hardware. At 1,000 RPS, use the simplest approach and buy bigger machines. At 1,000,000 RPS, the complexity of distributed systems is unavoidable and worth the maintenance cost because no single machine can handle that load.
The key insight is that complexity should be introduced deliberately, with documentation, tests, and team understanding. "We chose to add a custom cache eviction strategy because our access pattern is time-decaying and standard LRU wastes 40% of cache capacity" is a defensible complexity choice. "This code is complex because it evolved over 3 years without refactoring" is accidental complexity that harms both scalability and maintainability.
A common pattern that illustrates this tradeoff: teams often build custom distributed systems infrastructure (their own service mesh, their own configuration management, their own deployment pipeline) because off-the-shelf tools do not meet their exact needs. Each custom component is perfectly tuned for the current scale but costs months of engineering time to build and years to maintain. The alternative is accepting the 80%-fit of open-source tools (Kubernetes, Envoy, Terraform) and spending engineering time on the product instead of the platform. At Google-scale, custom infrastructure is justified. At most companies, it is premature complexity that slows everything down.
Reliability vs. Maintainability
Adding redundancy for reliability means more moving parts to maintain. Each failover mechanism, health check, retry loop, and circuit breaker is code that must be tested, monitored, and updated. A system with 5 layers of fault tolerance is more reliable but also harder to understand, debug, and modify than a simple system with no redundancy. When something goes wrong in a system with complex failover logic, the failover mechanism itself might be the source of the problem, and debugging a system that is actively trying to heal itself is notoriously difficult.
The balance here is to build reliability into the platform layer (use managed databases with automatic failover, deploy on Kubernetes with self-healing pods, use cloud load balancers with health checks) rather than reimplementing it in application code. Let the infrastructure handle fault tolerance so your application code stays simple and maintainable. A managed RDS instance with Multi-AZ failover gives you database reliability without writing or maintaining a single line of failover code.
There is a subtle trap here: teams sometimes add so many reliability layers that the interaction between them causes failures. A retry at the application layer combined with a retry at the load balancer layer combined with a retry at the client layer can turn a single failed request into 27 retry attempts (3 x 3 x 3), overwhelming the already-struggling service. This is called a "retry storm" and it is a case where reliability mechanisms actually reduce reliability. The lesson: reliability mechanisms must be coordinated across layers, not stacked blindly.
Prioritizing by Stage
The right balance shifts as your system matures, and recognizing which stage you are in prevents both premature optimization and under-investment.
Early-stage startups should prioritize maintainability above all else. Move fast, change direction easily, and accept lower scalability and reliability because you do not have millions of users yet. The biggest risk is not downtime or scale; it is building the wrong product. A simple, modifiable codebase lets you pivot quickly when you discover what users actually want.
Growth-stage companies should shift priority to scalability while preserving enough maintainability to keep shipping features. You have found product-market fit, traffic is growing, and your simple architecture is showing cracks. Invest in scaling the database, adding caching layers, and decomposing monolithic bottlenecks, but resist the urge to over-engineer for a scale you have not reached. This is also the stage where observability becomes critical: you need metrics and tracing to understand where the bottlenecks actually are, rather than guessing.
Mature systems at scale prioritize reliability because the cost of downtime is enormous. Millions of dollars per minute for large e-commerce platforms. Trust, once lost through an outage, takes months to rebuild. The system already has its users and its scale, so the marginal value of new features is lower than the marginal cost of an outage. At this stage, teams implement formal change management processes, extensive canary testing, and phased rollouts that would be overkill for a startup but are essential for a platform serving millions of users.
No system gets all three perfectly. The skill is knowing which one matters most right now, investing accordingly, and being willing to shift priorities as circumstances change.
A Decision Framework
When facing a design decision, ask three questions in order:
- What is our current stage? This determines which property gets priority. A pre-product-market-fit startup should almost always choose the option that maximizes development speed (maintainability). A mature platform should almost always choose the option that minimizes outage risk (reliability).
- What does the data say? Do not guess about load parameters, failure rates, or performance bottlenecks. Measure them. A team that "knows" their database is the bottleneck but has never profiled their queries is making decisions based on assumptions. Instrument first, then decide.
- What is the reversibility of this decision? Reversible decisions (choosing a cache TTL, picking an API design for an internal service) can be made quickly and cheaply corrected. Irreversible decisions (choosing a database technology, defining a data model, selecting a cloud provider) deserve more analysis because the cost of changing later is high.
This framework does not eliminate tradeoffs. It makes them explicit, data-driven, and appropriate to your context, which is the best any engineering team can do.
Common Anti-Patterns
Recognizing anti-patterns is as valuable as knowing best practices:
"We need microservices for scalability." Microservices solve organizational scaling (independent teams deploying independently) more than they solve technical scaling. A well-optimized monolith often outperforms a microservice architecture at the same scale because it avoids network serialization overhead between components. Microservices are worth their complexity when you have enough engineers that the monolith's deploy pipeline is the bottleneck, not the runtime performance.
"Let us build for a million users from day one." This leads to over-engineered systems that are expensive to operate and slow to modify, serving 100 users. Build for 10x your current scale, not 10,000x. When you reach 10x, you will understand your actual load parameters and can make informed scaling decisions instead of guessing.
"We cannot ship features until the codebase is clean." This is the maintainability trap: using code quality as a reason to delay shipping. The fix is incremental improvement, not perfection-before-progress. Ship the feature, clean up the code it touched, and repeat. Over time, the most-modified areas of the codebase improve naturally because they are touched most often.
"Our SLO is 99.99% because we care about quality." Choosing an SLO without calculating the cost is vanity engineering. Four nines means less than 4.3 minutes of downtime per month. Achieving this requires multi-region active-active deployments, automated failover with sub-minute detection, zero-downtime deploys, and comprehensive chaos testing. For an internal dashboard used by 50 employees, three nines (43 minutes per month) is more than sufficient, and the engineering resources saved can be invested in features that actually move the business forward.
"We should use the same architecture as Google/Netflix/Amazon." These companies operate at a scale that drives unique architectural decisions. Google built Bigtable because no existing database could handle their workload. Netflix built their own CDN because no vendor could match their streaming requirements. Unless you are operating at comparable scale, their architecture is a poor fit for your constraints. What you should learn from them is their decision-making framework (measure, then act), not their specific decisions.
Avoiding these anti-patterns requires intellectual honesty about your current situation. Measure your actual load parameters, not your projected ones. Profile your actual bottlenecks, not your assumed ones. Design for the system you have, not the system you wish you had. The best engineers are not the ones who build the most sophisticated systems. They are the ones who build the simplest system that solves the actual problem, and who know when the problem has changed enough to warrant a different solution.