English | 繁體中文
System design
System design is the practice of turning product behavior, traffic, data, and failure assumptions into an operable system. Start with the contract and constraints; choose the simplest architecture that satisfies them; make trade-offs explicit and measurable. A diagram is a hypothesis, not the design by itself.
A repeatable design workflow
- Clarify the scope: users, use cases, out-of-scope behavior, regions, tenancy, compliance, and freshness expectations.
- State assumptions: requests per second, object sizes, read/write ratio, peak multiplier, retention, growth, and failure budget.
- Write requirements: functional behavior plus measurable non-functional targets.
- Estimate capacity: storage, bandwidth, connections, CPU, memory, and cost; show the arithmetic.
- Define the data and API contracts: keys, invariants, pagination, errors, idempotency, and consistency.
- Draw the request path: clients, edge, services, caches, databases, queues, workers, and external dependencies.
- Walk through failure and growth: retries, overload, partial failure, recovery, migration, and the next bottleneck.
- Choose tools last: select managed or self-hosted components based on requirements, team skill, operational burden, and exit cost.
Rule of thumb: explain what happens on the happy path, then what happens when each dependency is slow, unavailable, duplicated, or inconsistent.
Example scenario and assumptions
Assume a multi-region URL-shortening service. Users create a short link and later redirect through it.
Traffic: 100 writes/s average, 1,000 reads/s average, 10x peak burst
Read/write: 10:1; redirect latency target p95 < 100 ms
Durability: created links must not be lost after acknowledgement
Availability: 99.95% monthly for redirects; writes may degrade before reads
Retention: 5 years; 1 KiB average record; 20% annual growth
Consistency: redirect must see a successful write within 5 secondsThese numbers are assumptions, not universal facts. Replace them before selecting capacity or a database.
Requirements
Functional requirements
- Create a short link with a destination, owner, expiration, and optional idempotency key.
- Redirect a short code to its current destination, returning a clear not-found or expired response.
- Revoke or update a link subject to authorization.
- List links by owner with stable, cursor-based pagination.
- Emit an audit event for create, update, revoke, and redirect policy decisions.
Non-functional requirements
| Attribute | Example target | Design consequence |
|---|---|---|
| Availability | Redirects 99.95% monthly | Replicas, health checks, graceful degradation |
| Latency | Redirect p95 < 100 ms | Edge/cache, bounded queries, no synchronous analytics |
| Durability | No acknowledged write lost | Replicated durable store, tested recovery |
| Consistency | Read-after-write within 5 s | Session routing, invalidation, or bounded replication lag |
| Scalability | 10x burst without data loss | Stateless tiers, autoscaling, queues, backpressure |
| Security | Tenant isolation and least privilege | Auth at the edge and service, scoped data access |
| Operability | Actionable alerts and rollback | Metrics, traces, logs, runbooks, versioned changes |
| Cost | Fit an explicit monthly budget | Managed services, tiered storage, right-sized capacity |
Do not say “fast,” “highly available,” or “scalable” without a target, workload, and measurement method.
Capacity estimation
Use orders of magnitude and preserve headroom:
peak RPS = average RPS × peak multiplier
annual records = writes/s × seconds/year
raw storage = records × average record size
bandwidth = RPS × average payload size
working set = hot-record fraction × total recordsFor the scenario, 100 writes/s × ~31.5M s/year ≈ 3.15B records/year. At 1 KiB each, raw data is roughly 3.15 TB/year before indexes, replicas, backups, and metadata. The estimate may invalidate the first storage choice; that is useful information. Include 2–3x operational headroom, not only the average.
Canonical request path
Client
→ DNS / CDN / WAF / rate limiter
→ load balancer
→ stateless API or redirect service
→ cache (hot short codes)
→ primary data store and read replicas
→ queue / event stream
→ asynchronous analytics, audit, and cleanup workers
→ object storage for exports, backups, and cold dataResponsibilities of common components
- DNS/CDN: route users, cache safe responses, absorb geography and bandwidth; do not cache personalized or revocable data accidentally.
- Load balancer: health-aware routing, TLS termination where appropriate, connection limits, and draining during deploys.
- Stateless service: authentication, validation, authorization, orchestration, and response shaping; keep durable state outside instances.
- Cache: reduce repeated reads; define TTL, invalidation, stampede protection, and stale-data policy before adding it.
- Primary store: source of truth for invariants and transactions.
- Replica/search index: optimize read patterns, but document lag and rebuild procedures.
- Queue/stream: decouple slow or bursty work; use durable delivery, bounded consumers, retry policy, and a dead-letter path.
- Object storage: inexpensive durable blobs, snapshots, exports, and replayable raw events.
- Observability stack: metrics for symptoms, logs for details, traces for request paths, and profiles for resource hotspots.
Data and API design
Choose access patterns before tables. Model ownership, uniqueness, lifecycle, and query keys explicitly.
POST /v1/links
Idempotency-Key: 9b2d...
Content-Type: application/json
{"destination":"https://example.com/docs","expiresAt":"2027-01-01T00:00:00Z"}{
"data": {
"id": "link_123",
"code": "aB91x",
"destination": "https://example.com/docs",
"expiresAt": "2027-01-01T00:00:00Z"
},
"meta": {"requestId": "req_456"},
"error": null
}- Validate syntax, size, allowed schemes, expiration, and tenant ownership at the boundary.
- Use stable opaque identifiers; never expose sequential IDs when they reveal sensitive volume or tenancy.
- Enforce uniqueness and state transitions in the database, not only in application code.
- Make retries safe with an idempotency key whose request fingerprint is stored with the result.
- Prefer cursor pagination over offset pagination for changing or large datasets.
- Version contracts deliberately; document error codes, retryability, and deprecation windows.
- Keep authorization predicates in every query scope, including cache keys and aggregations.
CAP theorem and BASE
CAP applies to a distributed data system during a network partition: it can guarantee at most two of Consistency, Availability, and Partition tolerance. Since partitions are possible in real distributed networks, the practical decision is usually how the system behaves during a partition:
- CP: preserve a single agreed value and reject or delay some requests. Choose for uniqueness, balances, locks, or state transitions where stale writes are harmful.
- AP: continue serving from reachable replicas and reconcile later. Choose for feeds, reactions, telemetry, or other workloads where availability and eventual convergence matter more than an immediate global value.
CAP does not mean “pick two forever,” nor does it say every endpoint has the same choice. A system can use CP behavior for link creation and AP-ish caching for redirects.
BASE—Basically Available, Soft state, Eventual consistency—is a useful AP-oriented approach:
- Basically available: return a useful response despite partial failure.
- Soft state: state may change without a new user request because replicas, caches, or background work converge.
- Eventual consistency: if writes stop and the system remains healthy, replicas eventually agree.
State the convergence bound and conflict rule. “Eventually” without a bound is not a requirement.
PACELC theorem
PACELC adds the normal operating case to CAP:
If there is a Partition, choose Availability or Consistency; Else, choose lower Latency or stronger Consistency.
Examples:
- A globally replicated database may favor low latency by reading a nearby replica, accepting replication lag.
- A leader-routed database may pay cross-region latency to provide fresher reads.
- A cache may favor latency normally, while invalidation or version checks protect correctness for selected operations.
Record the choice per operation, not as a slogan for the whole product.
Common design patterns
- Stateless horizontal scaling: instances are replaceable; session state belongs in a durable or shared store.
- Cache-aside: read cache, load on miss, then populate; define invalidation and stampede control.
- Write-through / write-behind: couple or decouple persistence from cache writes; write-behind needs durable buffering and explicit loss semantics.
- CQRS: separate write/invariant enforcement from read-optimized projections; accept projection lag and rebuild cost.
- Event-driven processing: publish durable domain events and process asynchronously; design for duplicate delivery and ordering scope.
- Outbox pattern: commit business state and an event record atomically, then publish the outbox asynchronously.
- Saga: coordinate a multi-step workflow with compensating actions when one distributed transaction is impractical.
- Bulkhead and circuit breaker: isolate dependency pools and stop multiplying failures; pair with timeouts and fallback behavior.
- Token bucket / leaky bucket: enforce rate or burst limits at a stated scope: IP, user, tenant, or API key.
- Sharding: partition by a stable, high-cardinality key; plan hot keys, resharding, cross-shard queries, and rebalancing first.
- Consistent hashing: reduce movement when nodes change; still monitor skew and hot partitions.
- Leader election / leases: coordinate ownership with expiry and fencing tokens; do not assume a lock alone prevents stale owners from acting.
Patterns are tools, not defaults. Each adds state, failure modes, and operational work.
Consistency, retries, and failure handling
- Set deadlines on every network hop; a retry must not outlive the user request or overload the dependency.
- Retry only transient, idempotent operations with bounded exponential backoff and jitter.
- Avoid retry storms: cap attempts, honor
Retry-After, use circuit breakers, and budget retries across a request tree. - Distinguish timeout, cancellation, validation failure, conflict, dependency failure, and unknown commit outcome.
- For queues, choose at-least-once delivery unless exactly-once semantics are proven end-to-end. Make consumers idempotent with a deduplication key.
- Use monotonic versions, compare-and-swap, or conditional writes for concurrent updates.
- Define ordering scope: global ordering is expensive; per-aggregate or per-partition ordering is often sufficient.
- Plan poison-message handling, dead letters, replay, backfill, and schema compatibility.
Tooling selection
| Need | Start with | Select based on |
|---|---|---|
| Relational invariants and transactions | PostgreSQL or managed equivalent | Joins, write contention, replication, operations |
| Key/value or very high-scale predictable access | Managed NoSQL store | Partition key quality, consistency modes, access patterns |
| Full-text or faceted search | OpenSearch/Elasticsearch or hosted search | Index freshness, relevance, shard operations, cost |
| Cache / ephemeral coordination | Redis-compatible service | Eviction, persistence, HA, memory cost, lock semantics |
| Durable asynchronous work | Managed queue or Kafka-compatible stream | Delivery, ordering, replay, throughput, retention |
| Blob and backup storage | Object storage | Lifecycle policy, retrieval cost, durability, egress |
| Service communication | HTTP/JSON first; gRPC for internal contracts | Debuggability, schema evolution, deadlines, tooling |
| Deployment | Containers plus managed orchestration | Team operations, autoscaling, portability, platform cost |
| Infrastructure | Terraform/OpenTofu or platform-native IaC | Reviewability, drift control, state security, maturity |
| Observability | OpenTelemetry + metrics/logs/traces backend | Cardinality, retention, alert quality, vendor lock-in |
Tool choice follows workload and team capability. A managed service often costs more per unit but less in engineering time, on-call load, patching, and recovery risk. A self-hosted component is justified only when its control, performance, compliance, or unit economics outweigh that burden.
Scalability and cost considerations
- Scale the bottleneck, not every component: measure CPU, memory, lock contention, database I/O, cache hit rate, queue age, and tail latency.
- Separate read and write paths when their scaling or consistency needs differ; do not split services merely by nouns.
- Cache intentionally: calculate hit rate, invalidation cost, memory footprint, and stale-data risk. A cache miss path must remain safe under a full cache outage.
- Protect dependencies: admission control, bounded queues, load shedding, quotas, and backpressure are often cheaper than overprovisioning.
- Use tiered storage: hot data in fast stores, cold data in object storage, and lifecycle policies for expiration and archival.
- Control egress and replication: cross-region traffic and duplicated storage can dominate compute cost.
- Estimate total cost of ownership: infrastructure, observability, backups, support plans, on-call time, migrations, and incident impact.
- Prefer reversible decisions: keep portable data formats, export paths, tested restore procedures, and clear ownership boundaries.
- Load-test the peak and failure modes: a healthy average benchmark does not reveal queue growth, cache stampedes, or retry amplification.
Security and operability checklist
- Authenticate and authorize every protected operation; enforce tenant isolation server-side.
- Validate and size-limit input; use parameterized queries and output encoding; protect SSRF, path traversal, and unsafe redirects.
- Encrypt in transit and at rest; keep secrets in a secret manager, rotate them, and never log them.
- Apply least privilege to services, operators, queues, buckets, and database roles.
- Record audit events without leaking credentials or sensitive payloads; define retention and access.
- Use structured logs with request/correlation IDs, redaction, and bounded cardinality.
- Alert on user-impacting symptoms: error rate, tail latency, saturation, queue age, replication lag, and failed backups.
- Test restore, failover, dependency outage, deploy rollback, schema migration, and clock/region failure—not just unit behavior.
- Define SLOs, error budgets, ownership, runbooks, and a communication plan before production.
Compilable Go example: bounded, cancellable worker stage
A queue-backed service often needs a bounded worker stage for asynchronous work. The example keeps ownership and shutdown explicit, propagates cancellation, and reports the first error without spawning unbounded goroutines. Save it as main.go and run go run ..
package main
import (
"context"
"errors"
"fmt"
"sync"
)
type Job struct {
ID int
Value int
}
type Result struct {
ID int
Square int
Err error
}
func worker(ctx context.Context, jobs <-chan Job, results chan<- Result, wg *sync.WaitGroup) {
defer wg.Done()
for {
select {
case <-ctx.Done():
return
case job, ok := <-jobs:
if !ok {
return
}
if job.Value < 0 {
results <- Result{ID: job.ID, Err: errors.New("negative values are invalid")}
continue
}
results <- Result{ID: job.ID, Square: job.Value * job.Value}
}
}
}
func run(ctx context.Context, jobs []Job, workerCount int) ([]Result, error) {
if workerCount < 1 {
return nil, errors.New("worker count must be positive")
}
ctx, cancel := context.WithCancel(ctx)
defer cancel()
jobCh := make(chan Job)
resultCh := make(chan Result)
var workers sync.WaitGroup
workers.Add(workerCount)
for i := 0; i < workerCount; i++ {
go worker(ctx, jobCh, resultCh, &workers)
}
go func() {
defer close(jobCh)
for _, job := range jobs {
select {
case jobCh <- job:
case <-ctx.Done():
return
}
}
}()
go func() {
workers.Wait()
close(resultCh)
}()
results := make([]Result, 0, len(jobs))
var firstErr error
for result := range resultCh {
results = append(results, result)
if result.Err != nil && firstErr == nil {
firstErr = fmt.Errorf("job %d: %w", result.ID, result.Err)
cancel()
}
}
return results, firstErr
}
func main() {
results, err := run(context.Background(), []Job{{1, 3}, {2, 4}, {3, -1}}, 2)
fmt.Printf("results=%v error=%v\n", results, err)
}This is a building block, not a complete queue system. Production code still needs durable enqueue semantics, idempotent processing, metrics, retry/backoff policy, dead-letter handling, and graceful process shutdown.
Final design review questions
- What are the exact read/write paths, keys, invariants, and consistency expectations?
- Which data is authoritative, and how are replicas, caches, and projections repaired?
- What happens during a partition, full cache, slow dependency, duplicate message, or unknown write outcome?
- What is the first bottleneck at 1x, 10x, and 100x traffic?
- Can the team operate, restore, migrate, and debug every chosen component?
- Which trade-off was chosen—latency, consistency, availability, durability, simplicity, or cost—and how will it be measured?
The best design is not the one with the most components. It is the smallest system whose behavior, failure modes, and operating cost are understood well enough to evolve safely.