Learn Labs
1. Trade-Offs in Data Systems Architecture

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

ReasonExplanation
Inherent distributionTwo or more interacting users on their own devices — communication must cross a network
Requests between cloud servicesData stored in one service, processed in another → network transfer. Cloud native & microservices are therefore distributed
Fault tolerance / HARedundancy so that a machine, several machines, the network, or a whole datacenter can fail and another takes over
ScalabilityData volume or compute exceeds one machine
LatencyServers in multiple regions so each user is served from geographically nearby
ElasticityScale up/down with demand and pay only for what you use, instead of provisioning for peak
Specialized hardwareObject store = many disks/few CPUs; analytics = lots of CPU+RAM/no disks; ML = GPUs
Legal complianceData residency laws requiring data about people in a jurisdiction to be stored/processed inside it (scope varies — sometimes only medical or financial data)
SustainabilityRun 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.

DimensionSupercomputing / HPCCloud computing
WorkloadComputationally intensive science: weather forecasting, climate modeling, molecular dynamics, optimization, PDEsOnline services, business data systems serving user requests with high availability
Fault handlingLarge batch jobs checkpoint to disk; on node failure, stop the whole cluster, repair, restart from last checkpointStopping the cluster is unacceptable — must serve users continuously with minimal interruption
CommunicationShared memory and RDMA — high bandwidth, low latency, assumes high trust among usersNetwork and machines shared by mutually untrusting organizations → VMs for resource isolation, encryption, authentication
Network topologySpecialized: multidimensional meshes and toruses, tuned to known communication patternsIP + Ethernet in Clos topologies for high bisection bandwidth
GeographyAll nodes assumed close togetherNodes distributed across geographic regions

On this page