Microservices and Distributed Systems Basics
"Should this be a microservice?" and "what goes wrong when it is distributed?" are staple backend questions. A strong answer is balanced: microservices solve organisational and scaling problems, but they trade simple function calls for networks, partial failure and eventual consistency. This chapter covers when to split, how to draw boundaries, how services communicate, and the distributed-systems realities (time, consistency, consensus, idempotency) you must reason about.
1. Monolith, modular monolith, microservices
| Style | Description | Strengths | Costs |
|---|---|---|---|
| Monolith | one deployable unit | simple development, debugging, transactions and deployment; fast in-process calls | growing coupling, a single scaling unit, slower builds and deploys, a shared blast radius |
| Modular monolith | one deployable with enforced internal module boundaries | most of the monolith's simplicity plus clear seams for later extraction | needs discipline to keep boundaries clean |
| Microservices | many small, independently deployable services, each owning its data | independent deploys and scaling, team autonomy, technology choice, fault isolation | distributed-systems complexity, operations, consistency, testing and debugging difficulty |
Start with a (modular) monolith unless you have clear reasons. Split when: separate teams are blocked on each other's releases, components have very different scaling or reliability needs, or a domain is genuinely independent. Splitting too early produces a distributed monolith: all the costs of microservices with none of the independence (services that must be deployed together, share a database or call each other in long synchronous chains).
2. Drawing service boundaries
- By business capability / bounded context (domain-driven design): orders, payments, inventory, catalogue, notifications, each with its own model and vocabulary. "Customer" means different things in billing and in support, so each context keeps its own representation.
- High cohesion, low coupling: things that change together live together; the interface between services is narrow and stable.
- Data ownership: each service owns its data store and exposes it only through its API or events. Shared databases couple services through the schema and defeat independent deployment.
- Team alignment (Conway's law): the architecture tends to mirror the communication structure of the organisation. Align boundaries with teams that own services end to end ("you build it, you run it").
- Right size: not a tiny function ("nano-service"), not an entire product. A team of a few engineers should be able to understand and own it.
3. Communication styles
| Synchronous (request-response) | Asynchronous (messaging) | |
|---|---|---|
| Examples | REST, gRPC | queues, topics, event streams |
| Coupling | temporal: caller needs the callee to be up | looser: the broker buffers |
| Latency | caller waits | caller continues |
| Failure | propagates to the caller (needs timeouts, retries, breakers) | handled by retries and dead-letter queues |
| Fits | queries, commands needing an immediate answer | notifications, workflows, integration, load levelling |
Long chains of synchronous calls multiply latency and multiply failure probability: if five services each have 99.9% availability and must all respond, the overall availability is about .
def chain_availability(availabilities):
result = 1.0
for a in availabilities:
result *= a
return result
assert abs(chain_availability([0.999] * 5) - 0.995) < 1e-4
assert chain_availability([0.999] * 20) < 0.981 # twenty hops: the system is down about 2 % of the time
def parallel_redundancy(a, copies):
return 1 - (1 - a) ** copies # any one healthy replica is enough
assert parallel_redundancy(0.99, 2) > 0.9998
API gateway and backend for frontend
An API gateway is the single entry point: routing, authentication, rate limiting, TLS termination, request aggregation, protocol translation. A backend for frontend (BFF) is a gateway tailored to one client type (web, mobile) that shapes responses for it. Keep business logic out of the gateway.
Service discovery and load balancing
Instances come and go, so clients find them through service discovery (DNS, Consul, Kubernetes services, a service mesh). Load balancing can be server-side (a proxy), or client-side (the caller picks an instance); algorithms include round-robin, least-connections, and power of two choices (pick two random instances, use the less loaded one), which avoids herding.
import random
def least_loaded_of_two(loads, rng):
a, b = rng.sample(range(len(loads)), 2)
return a if loads[a] <= loads[b] else b
def simulate(strategy, servers=50, requests=5000, seed=1):
rng = random.Random(seed)
loads = [0] * servers
for _ in range(requests):
i = rng.randrange(servers) if strategy == "random" else least_loaded_of_two(loads, rng)
loads[i] += 1
return max(loads) - min(loads)
assert simulate("two-choices") < simulate("random") # a tiny bit of load information balances far better than pure chance
Service mesh
A service mesh (Istio, Linkerd) moves cross-cutting network concerns (mutual TLS, retries, timeouts, circuit breaking, traffic splitting, telemetry) into sidecar proxies managed by a control plane, so each service need not implement them. It adds operational weight; adopt it when you have many services and a platform team.
4. Data in a microservice world
- Database per service. Queries that used to be joins become API calls or replicated read models built from events.
- No distributed transactions in practice (two-phase commit is slow and fragile across services). Use sagas with compensations and the outbox pattern (see the messaging chapter).
- Eventual consistency is the default; design the user experience around it (pending states, "we are processing your order").
- Data duplication is acceptable when it is derived from the owner's events and treated as a cache; the owner remains the source of truth.
- Reporting and analytics: stream events or CDC to a warehouse; do not run analytical queries against operational databases.
5. The realities of distributed systems
The eight fallacies: the network is reliable; latency is zero; bandwidth is infinite; the network is secure; topology does not change; there is one administrator; transport cost is zero; the network is homogeneous. Each is false.
Partial failure and the two generals problem
A caller that gets no response does not know whether the request was lost, the server crashed before acting, or it acted and the reply was lost. This uncertainty is why timeouts need idempotent retries, and why exactly-once delivery across a network cannot be guaranteed.
Time and ordering
Clocks on different machines drift, so you cannot order events across machines by wall-clock time. Tools:
- Logical clocks: a Lamport clock gives a counter consistent with causality (if A happened before B, then ).
- Vector clocks: detect whether two events are causally ordered or concurrent.
- Hybrid logical clocks combine physical and logical time.
- TrueTime-style bounded uncertainty (Spanner) uses specialised hardware.
- Sequence numbers from a single log give a total order within a partition.
class LamportClock:
def __init__(self):
self.t = 0
def tick(self): # a local event
self.t += 1
return self.t
def send(self):
return self.tick() # the message carries the sender's timestamp
def receive(self, msg_t):
self.t = max(self.t, msg_t) + 1 # jump ahead of anything the sender had seen
return self.t
a, b = LamportClock(), LamportClock()
a.tick(); a.tick()
ts = a.send() # a's third event
b.tick()
recv = b.receive(ts)
assert recv > ts # the receive is ordered after the send, whatever the wall clocks say
def vc_compare(x, y):
le = all(x[k] <= y.get(k, 0) for k in x); ge = all(y[k] <= x.get(k, 0) for k in y)
return "equal" if le and ge else "before" if le else "after" if ge else "concurrent"
assert vc_compare({"A": 2, "B": 0}, {"A": 2, "B": 1}) == "before"
assert vc_compare({"A": 3, "B": 0}, {"A": 2, "B": 1}) == "concurrent" # neither saw the other: a true conflict
Consistency models and CAP
(See the schema and NoSQL chapter for CAP, PACELC and quorums.) For user-facing behaviour, name the guarantee you need: read-your-writes (a user sees their own update), monotonic reads (never go back in time), causal consistency, or linearizability (behaves like a single copy). Stronger guarantees cost latency and availability.
Replication and conflicts
With leaderless or multi-leader replication, concurrent writes to the same item conflict. Strategies: last-write-wins (simple, loses data), version vectors with application-level merge, and CRDTs (data types designed so replicas always converge, such as counters and sets).
class GCounter:
"""A grow-only counter CRDT: each replica increments its own slot; merge takes the maximum per slot."""
def __init__(self, replica):
self.replica, self.counts = replica, {}
def increment(self, n=1):
self.counts[self.replica] = self.counts.get(self.replica, 0) + n
def merge(self, other):
for r, c in other.counts.items():
self.counts[r] = max(self.counts.get(r, 0), c)
def value(self):
return sum(self.counts.values())
x, y = GCounter("x"), GCounter("y")
x.increment(3); y.increment(2) # concurrent updates on two replicas, with no coordination
x.merge(y); y.merge(x)
assert x.value() == y.value() == 5 # both converge to the same total
x.merge(y); x.merge(y)
assert x.value() == 5 # merging again is harmless (idempotent, commutative, associative)
Consensus and leader election
Replicas that must agree (who is the leader, what is the next log entry) use consensus protocols such as Raft and Paxos, which tolerate failures with nodes by requiring a majority quorum. You rarely implement them; you use systems built on them (etcd, ZooKeeper, Consul, Spanner, CockroachDB, Kafka's KRaft). Know what a leader election and a quorum are, and why an even cluster size adds cost without adding fault tolerance.
def fault_tolerance(nodes):
quorum = nodes // 2 + 1
return nodes - quorum # how many nodes can fail while a majority remains
assert [fault_tolerance(n) for n in (3, 4, 5, 6, 7)] == [1, 1, 2, 2, 3] # 4 nodes tolerate no more failures than 3
Split brain
If a partition lets two nodes both believe they are the leader, they can accept conflicting writes. Defences: majority quorums, fencing tokens and leases with proper expiry, and STONITH (shoot the other node in the head) in clustering setups.
6. Distributed tracing and debugging
A request touches many services, so logs from one service are not enough. Propagate a correlation (trace) ID across calls and messages, record spans for each hop, and view the whole path (see the observability chapter). Without it, debugging a latency problem across ten services is guesswork.
7. Deployment and operations of many services
- Containers and orchestration (Docker, Kubernetes) for packaging, scaling and self-healing.
- CI/CD per service with automated tests; contract tests so a provider cannot break its consumers unnoticed.
- Independent versioning and backward-compatible APIs; use expand-and-contract for breaking changes.
- Configuration and secrets managed centrally; feature flags.
- Observability (logs, metrics, traces) and SLOs per service.
- Platform engineering: a paved road of templates, shared tooling and standards, so each team need not reinvent them.
8. Migration: the strangler fig pattern
To move from a monolith to services without a big-bang rewrite: put a facade (gateway) in front, extract one capability at a time behind it, route traffic gradually, and retire the old code path when the new one is proven. Extract along seams that are already loosely coupled, and give each new service its own data (migrating it carefully, often with dual writes or CDC during the transition).
9. Trade-off questions
| Question | Reasoning |
|---|---|
| Microservices or monolith for a 5-person startup? | monolith: fewer moving parts, faster iteration, split later along clear seams |
| Sync or async between services? | sync when the caller needs an answer now; async for side effects, integration and resilience |
| Shared database or database per service? | per service for independence; a shared database couples release cycles and schema |
| Strong or eventual consistency? | strong for money and inventory at the point of commitment; eventual for feeds, search, analytics |
| Orchestration or choreography? | orchestration for visibility in complex flows; choreography for simple, loosely coupled reactions |
| Build or buy (managed services)? | buy for undifferentiated infrastructure; build where it is your core |
10. Common mistakes
- Starting with microservices for a small system, or splitting by technical layer (a "database service", a "UI service") instead of by business capability.
- A distributed monolith: shared databases, lock-step deployments, chatty synchronous calls.
- Ignoring partial failure: no timeouts, retries or fallbacks.
- Distributed transactions by default instead of sagas and idempotency.
- No contract testing or versioning, breaking consumers.
- Treating wall-clock time as a reliable global order.
- Underestimating operational cost: monitoring, deployment pipelines, on-call load.
- Too many services with too few engineers.
11. Practice questions
- When would you choose microservices over a monolith, and when not?
- How do you decide where to cut a monolith? What is a bounded context?
- How do you handle a transaction that spans three services?
- Explain why a chain of synchronous calls hurts availability and latency, and how to mitigate it.
- What is the CAP theorem? What does a quorum of guarantee?
- Why can't you rely on timestamps to order events across servers? What can you use?
- What is consensus, and why does a 5-node cluster survive two failures?
- Describe the strangler fig pattern for migrating off a monolith.