PrepZone Logo
PrepZone

What Makes a System Distributed

Independent nodes, partial failure, and no shared clock — the three facts every distributed design must accept.

Why this matters

  • Every backend interview eventually probes partial failure — you must reason about nodes, not a single machine.
  • Production outages come from assuming one process, one clock, and one reliable network.
  • ShardPay starts as one JVM and grows into replicated shards — this article frames every later chapter.
  • The gap between "it works on my laptop" and "it works at 2 a.m. during failover" is almost always a distribution problem.
ClientMobile / web
API gateway
Ledger AShard 1
Ledger BShard 2
Ledger CShard 3
Independent nodes communicate over a network. Any node can fail without warning.

Independent nodes over an unreliable network

ShardPay's ledger service runs on three JVMs in two availability zones. A transfer request hits any API instance, but each ledger node holds only part of the account data. No node sees the full cluster state at once — that is distribution. The mobile app still sees one balance and one transfer button; behind the gateway, requests fan out across shards, health checks, and retry policies.

Distribution is not just "many servers." It is independent failure domains connected by a network that drops packets, adds latency, and reorders messages. A developer who treats ledgerClient.transfer() like a local method call will eventually ship a bug that only appears when one AZ is slow and another is healthy.

Key points

  • Node — a process on a machine that communicates with peers over the network rather than shared memory. At ShardPay, each ledger JVM is a node with its own heap, disk, and election state; a crash in node B does not stop node A from serving reads for accounts it owns.
  • Partial failure — some components work while others are slow, partitioned, or dead; there is no clean global crash like a single-process OOM. ShardPay once returned HTTP 200 with a stale balance because the follower replica was reachable but had not caught up after leader failover.
  • Shared-nothing — each node owns its memory and disk; coordination happens through explicit protocols (RPC, consensus, message queues), not a single shared database session. ShardPay shards partition accounts by accountId hash so no node needs the full account table in RAM.
  • Transparency — users and product APIs see one coherent system; engineers must see many moving parts with different consistency guarantees. The public POST /v1/transfers endpoint hides shard routing, idempotency keys, and cross-shard saga orchestration.

Why distribution is hard

On a single machine, a crash stops everything — easy to detect, easy to reason about. In a cluster, a node can hang without responding, return garbage after a long GC pause, or deliver a success response after the client already timed out and retried. The hardest failures are ambiguous: you cannot tell if the request was lost, still in flight, or completed on a node you no longer trust.

ShardPay learned this during its first multi-AZ deployment. A follower served a balance from 30 seconds ago while the leader had already committed a debit. The HTTP response was fast and well-formed — the error was silent staleness, not a 500. Distributed systems force you to design for uncertainty: timeouts, version vectors, fencing tokens, and explicit consistency modes.

Walkthrough: one transfer across three nodes

A customer transfers $50 from account A-1042 (shard 1) to B-8831 (shard 2). The API gateway picks a healthy ledger coordinator. The coordinator debits shard 1, then calls shard 2 to credit. If the network partitions between those two RPCs, shard 1 may hold a debit with no matching credit — a classic distributed transaction problem that module 10 addresses with sagas. The point here: three independent steps, three failure points, one user expectation of atomicity.

Java
// ShardPay — ledger node interface; each JVM implements this contract
public interface LedgerNode {
    /** Routed to the shard that owns fromAccount. May involve cross-shard RPC. */
    TransferResult transfer(TransferRequest request);

    /** Strong read from leader, or bounded-staleness read from follower. */
    Balance getBalance(AccountId account, ReadConsistency consistency);
}

public record TransferRequest(
    String idempotencyKey,
    AccountId fromAccount,
    AccountId toAccount,
    long amountCents
) {}

Concurrency without a global clock

There is no single "now" in a distributed system. Wall clocks drift; NTP corrections jump backward; a "slow" node may process events in a different order than a "fast" one. ShardPay never uses System.currentTimeMillis() to order ledger entries — it uses monotonic sequence numbers assigned by the shard leader after consensus.

This matters for debugging too. A log line at 12:00:01.003 on node A does not prove it happened before 12:00:01.002 on node B. Distributed tracing (module 11) and logical clocks (module 7) exist because time is local, not global.

Ordering and visibility

  • Happens-before — if event A causally leads to B, every correct node must eventually observe A before B. ShardPay's transfer ID is generated only after the debit is durable on the leader, so downstream fraud checks always see a committed debit.
  • Eventual visibility — replicas converge, but not instantly. ShardPay's mobile app reads from a follower with a 500ms max staleness SLA for balance display; the transfer confirmation screen reads from the leader.
  • Split brain — two nodes both believe they are leader and accept conflicting writes. ShardPay's Raft layer rejects writes from stale leaders using term numbers and fencing tokens.

When distribution is worth the cost

Not every service needs a cluster on day one. ShardPay ran as a single JVM for its first year, handling 2,000 transfers/sec on one machine. Distribution entered when proven limits appeared: disk I/O on the ledger, independent scaling for fraud GPU workloads, and regulatory requirements for multi-AZ durability.

The goal is not "microservices for microservices' sake." The goal is matching architecture to constraints: throughput, fault isolation, team autonomy, and geographic presence. Every chapter in this track assumes you can articulate why ShardPay distributed a particular component, not just how.

Quick recall

Everything you need if you only revisit this box.

  1. Distribution means independent nodes and partial failure — not just many servers.
  2. Users see one system; engineers must design for ambiguous failures and local time.
  3. ShardPay evolves from one JVM to sharded replicas; every later module builds on this frame.

Test yourself

Answer these before moving on — recall is what makes it stick.