1.3 Distributed vs Single-Node Systems
A distributed system = several machines communicating over a network.
A distributed system = several machines communicating over a network. Each participating process is a node.
3.1 Nine legitimate reasons to go distributed
| Reason | Explanation |
|---|---|
| Inherent distribution | Two or more interacting users on their own devices — communication must cross a network |
| Requests between cloud services | Data stored in one service, processed in another → network transfer. Cloud native & microservices are therefore distributed |
| Fault tolerance / HA | Redundancy so that a machine, several machines, the network, or a whole datacenter can fail and another takes over |
| Scalability | Data volume or compute exceeds one machine |
| Latency | Servers in multiple regions so each user is served from geographically nearby |
| Elasticity | Scale up/down with demand and pay only for what you use, instead of provisioning for peak |
| Specialized hardware | Object store = many disks/few CPUs; analytics = lots of CPU+RAM/no disks; ML = GPUs |
| Legal compliance | Data residency laws requiring data about people in a jurisdiction to be stored/processed inside it (scope varies — sometimes only medical or financial data) |
| Sustainability | Run jobs where/when renewable electricity is plentiful and the grid isn't strained; cuts carbon and cost |
3.2 The costs, stated bluntly
- Every network request may fail or time out. When a request times out, you don't know whether the service received it, so retrying may not be safe. (Ch 9.)
- Network calls are vastly slower than in-process function calls, even in fast datacenter networks. With large data volumes it's often faster to bring the computation to the data than to move the data to a separate machine.
- More nodes are not always faster — a simple single-threaded program on one computer can significantly beat a cluster of 100+ CPU cores.
- Troubleshooting is hard. If the system is slow, where is the problem? This is the domain of observability: collecting data about execution and querying it so both high-level metrics and individual events can be analyzed. Tracing: OpenTelemetry, Zipkin, Jaeger — track which client called which server for which operation and how long it took.
- Cross-service consistency becomes the application's problem when each service owns its database. Distributed transactions (Ch 8) exist but are rarely used with microservices: they cut against service independence, and many databases don't support them.
Therefore: doing a task on a single machine is often much simpler and cheaper. CPUs, memory, and disks have grown larger, faster, more reliable. With single-node engines like DuckDB, SQLite, KùzuDB, many workloads now fit on one node.
Heuristic: don't rush into distribution. Prove you need it.
3.3 Microservices
Client/server over HTTP is the common distribution style; the same process is often both a server (handling requests) and a client (making outbound ones). SOA → refined into microservices: a service has one well-defined purpose (S3's is file storage), exposes a network API, and has one team responsible for it.
Advantages: independent updates → less cross-team coordination; per-service hardware allocation; implementation hidden behind an API so owners can change internals freely. Each service typically owns its own database — sharing a database would make the entire schema part of the service's API (undeployable, unchangeable) and let one service's queries hurt another's performance.
Costs:
- Testing requires running all dependencies
- Each service needs infra for releases, resource scaling, log collection, health monitoring, on-call alerting — which is why Kubernetes became the standard foundation
- API evolution is hard. Clients expect certain fields; adding/removing them breaks clients, and the breakage is often discovered late, in staging or production. OpenAPI and gRPC help manage the client/server API relationship (Ch 5).
Microservices are primarily a technical solution to a people problem — letting teams progress independently without coordinating. Valuable at a large company; in a small company with few teams it's likely unnecessary overhead, and the simplest possible implementation is preferable.
3.4 Serverless / FaaS
The cloud provider automatically allocates and frees hardware based on incoming requests — no explicit start/stop of instances. Just as cloud storage replaced capacity planning with metered billing, serverless brings metered billing to code execution: pay for the time your code runs.
Costs and caveats:
- Time limits on function execution, restricted runtime environments
- Cold starts on first invocation
- "Serverless" is a misleading name — each execution still runs on a server, just possibly a different one each time
- The term has been stretched: BigQuery and various Kafka offerings use "serverless" to mean autoscaling + billing by usage rather than by machine instance
3.5 Cloud computing vs supercomputing (HPC)
Useful because large-scale analytical systems sometimes look like HPC.
| Dimension | Supercomputing / HPC | Cloud computing |
|---|---|---|
| Workload | Computationally intensive science: weather forecasting, climate modeling, molecular dynamics, optimization, PDEs | Online services, business data systems serving user requests with high availability |
| Fault handling | Large batch jobs checkpoint to disk; on node failure, stop the whole cluster, repair, restart from last checkpoint | Stopping the cluster is unacceptable — must serve users continuously with minimal interruption |
| Communication | Shared memory and RDMA — high bandwidth, low latency, assumes high trust among users | Network and machines shared by mutually untrusting organizations → VMs for resource isolation, encryption, authentication |
| Network topology | Specialized: multidimensional meshes and toruses, tuned to known communication patterns | IP + Ethernet in Clos topologies for high bisection bandwidth |
| Geography | All nodes assumed close together | Nodes distributed across geographic regions |