Special offer. See more information about Outbyte and uninstall instructions. Please review EULA and Privacy policy.

What is a distributed system? It is a group of independent computers that coordinate over a network to provide a service or manage shared data. The central challenge is that parts of the system can fail or become slow independently: one machine may stop, or messages between healthy machines may be delayed or lost. Understanding that partial failure explains why distributed systems use replication, consensus, and carefully chosen trade-offs.

What is a distributed system?

A distributed system consists of separate computers—often called nodes—that communicate over a network and work together. To a user, the computers may look like one service, such as a database or application. Underneath, however, each node has its own processor, memory, storage, and view of what is happening.

That independence is useful: work can be divided across machines, and a service can continue when a component fails. It also creates the defining difficulty. A node cannot always tell whether another node has crashed, is merely slow, or is isolated by a broken network path. A request or reply may be delayed or lost, even when the machines themselves are running.

This is called partial failure. In a single computer, a process can often rely on local communication and a common clock. Across a network, participants must make decisions with incomplete information. A system therefore needs rules for such cases: whether to wait, retry, reject a request, elect a new leader, or allow a response based on data that may not be current.

Special offer. See more information about Outbyte and uninstall instructions. Please review EULA and Privacy policy.

How does the CAP theorem work?

CAP describes a trade-off that becomes important when a network partition prevents some nodes from communicating. In AWS’s explanation, consistency means a read sees the latest write or returns an error; availability means each request receives a non-error response; and partition tolerance means the system continues operating despite lost messages between nodes. AWS’s CAP theorem overview defines these terms.

During a partition, a system cannot guarantee both that every successful read reflects the latest write and that every request receives a successful response. It must decide how to handle operations whose correctness depends on information it cannot obtain:

  • Favor consistency: reject or delay an operation when the system cannot confirm it is safe. Some requests may fail, but the system avoids claiming success with uncertain or stale state.
  • Favor availability: respond even when isolated nodes cannot coordinate. Requests can continue, but a response may use stale data, or separate parts of the system may temporarily accept conflicting updates.

This is not a simple label for a whole product that tells you it is always “CP” or “AP.” The relevant question is what the design does for particular operations when a partition occurs. Networked systems generally need to account for partitions; the design choice is how to behave during them. CAP is specifically about partition behavior, not a general ranking of systems or a claim that consistency is always preferable.

What CAP does not say about normal operation

CAP focuses on partitions. When the network is working normally, there can still be a trade-off between consistency and latency: waiting for coordination can make an operation safer but slower. PACELC extends the discussion to include this normal-operation choice alongside the partition-time trade-off. AWS Builders’ Library discusses CAP and PACELC.

Special offer. See more information about Outbyte and uninstall instructions. Please review EULA and Privacy policy.

What is the difference between replication and consensus?

Replication means keeping multiple copies of data or service state on different nodes. Redundant copies can help a service stay available if a node fails, but replicas must be coordinated: they need rules for which updates count, how copies catch up, and what a read may return.

Consensus is a way for distributed participants to agree on a value or decision despite certain failures. It can be used to choose a leader, decide whether a queue entry is committed, or agree on a datastore value, as described by Google’s Site Reliability Engineering material. Consensus is not the same thing as copying data; it can help replicas agree on the order or outcome of changes that they apply.

For example, a replicated service might have several nodes that can store copies, but it still needs a protocol for deciding which update is authoritative if two nodes cannot communicate. Consensus protocols commonly use a quorum—a sufficient number of participants—to make decisions. A majority quorum prevents two competing majorities from independently committing conflicting decisions under the protocol’s assumptions.

How many nodes do I need for fault tolerance?

There is no universal node count. It depends on the failure model, the quorum rule, and whether the system is designed for crash failures or Byzantine behavior. For majority-based crash-fault tolerance, Google SRE (2017) gives the relationship 2f + 1 replicas to tolerate f crash failures. It gives 3f + 1 replicas for f Byzantine failures, where faulty participants may behave arbitrarily rather than simply stop. These are protocol design relationships, not a recommendation that every application should use a particular cluster size.

Special offer. See more information about Outbyte and uninstall instructions. Please review EULA and Privacy policy.
Failure model Replica relationship What the relationship means
Crash failure 2f + 1 replicas May tolerate f crashed replicas when a majority is required. Google SRE, 2017.
Byzantine failure 3f + 1 replicas May tolerate f Byzantine faulty members. Google SRE, 2017.

With three replicas and a majority requirement, the crash-failure relationship permits one replica to fail while the remaining two still form a majority. That arithmetic does not guarantee that a deployed service will remain usable: replicas need independent enough infrastructure, the application must handle failures correctly, and a majority must still be able to communicate.

Fault tolerance means maintaining availability through redundancy when another subsystem takes over work after a failure; that is AWS’s definition of the goal. AWS’s fault-tolerance overview explains the concept. Simply adding nodes is not enough if they share a power source, network path, or other failure point. The design must consider the failures it is intended to survive.

How should you compare distributed-system designs?

There is no single best consensus algorithm or architecture for every workload. Google SRE (2017) states: “There is no one ‘best’ distributed consensus and state machine replication algorithm for performance, because performance is dependent on a number of factors relating to workload, the system’s performance objectives, and how the system is to be deployed.” Compare designs against the needs and limits of the service rather than choosing by name alone.

  • Consistency: What must a read guarantee about prior writes? Can the application accept stale results or temporarily divergent copies?
  • Partition behavior: Which operations should fail, wait, or continue if nodes cannot communicate?
  • Failure model and quorum: Are you handling crashed nodes, arbitrary faulty members, or both? What number of participants is required to make progress?
  • Leadership and ordering: How does the system select a leader, and how does it establish which updates come first?
  • Performance: What latency and throughput does the workload require, including the delay introduced by coordination?
  • Operational burden and cost: Can the team deploy, monitor, recover, and maintain the design at the scale and complexity it introduces?

Use CAP to reason about what happens during a partition. Use PACELC when evaluating the additional latency-versus-consistency choice during normal operation. Neither framework picks an architecture for you; they help make the consequences of a choice explicit.

Special offer. See more information about Outbyte and uninstall instructions. Please review EULA and Privacy policy.
Independent reader supportYour contribution helps us test, update, and keep practical guides available for everyone.Support on Ko-Fi

How do I learn distributed systems with Kubernetes?

Kubernetes offers a practical way to connect the concepts to real deployments. Its official tutorials include an interactive basics path as well as examples involving Redis configuration, StatefulSets, Cassandra, and ZooKeeper. Start with the basics, then use the stateful examples to examine why applications need stable identity, persistent data, or coordinated state. The Kubernetes tutorials are the official starting point.

Kubernetes documentation also describes production control planes distributed across multiple computers and clusters with multiple nodes for fault tolerance and high availability. Its high-availability topology guidance explains control-plane layouts. For workloads spanning locations, the multiple-zone guidance treats regions, zones, and nodes as fault domains and recommends topology controls to spread workloads.

A learning sequence

  1. Build the vocabulary: Work through the Kubernetes interactive basics tutorial. Identify the nodes and components involved rather than treating the cluster as one machine.
  2. Study stateful examples: Follow the Redis configuration and StatefulSet tutorials, then examine the Cassandra and ZooKeeper examples. Ask which information must survive a restart and which components need to coordinate.
  3. Draw the failure domains: Map nodes, zones, and regions in the deployment guidance. Note which components would be affected by losing one node or one communication path.
  4. Reason through failures: For each example, ask what a user sees if a node becomes unavailable, if a request is retried, or if two parts of the system cannot communicate. Trace whether the system waits, rejects work, elects a leader, or risks serving stale state.
  5. Observe safely in a small cluster: Build a small learning cluster and observe leader changes and retries as you explore node and network unavailability. Treat this as a way to test your understanding, not proof that a production service will behave the same way.

The goal is not merely to memorize CAP or quorum formulas. It is to be able to explain what state the system trusts, how it decides, and what the user experiences when communication or a component fails.

Product prices and availability are accurate as of the date/time indicated and are subject to change. Any price and availability information displayed on Amazon at the time of purchase will apply.

Special offer. See more information about Outbyte and uninstall instructions. Please review EULA and Privacy policy.