Parallel and Distributed Computing
A structured guide to designing parallel and distributed systems through decomposition, communication, coordination, state management, scalability, and fault tolerance.
Foundations: Parallel and Distributed Systems
Parallel and distributed systems both seek to divide work, but they operate under different conditions. executes independent work simultaneously, often within one machine. divides computation and state among independent machines connected by a network.
A distributed system can be viewed as a collection of processes that exchange messages. Network communication takes time, processors run at different speeds, and each process has only a limited view of the whole system. Consequently, events cannot always be placed into one universally visible chronological order. Causal ordering is more useful: if one event could have influenced another, the system should represent that relationship.
The central design trade-off is that additional machines can provide more capacity and redundancy, but they also introduce more communication paths, possible event orderings, and independent failure modes.
Takeaway: Parallel systems primarily manage simultaneous execution; distributed systems must manage simultaneous execution together with communication uncertainty and partial knowledge.
and Work Distribution
is the first major design decision. A problem should be divided into portions that can make progress independently, while data placement and dependencies should be made explicit.
Three common forms are:
Task : different workers perform different operations, such as reading input, transforming it, and storing results.
Data : workers perform the same operation on different partitions of the data.
Pipeline : different stages process different inputs at the same time, so one input can be transformed while another is being read.
A useful example is large-scale word counting. Documents are partitioned among workers. Each worker emits intermediate pairs such as (word, 1). The system then groups pairs by word and adds the associated values. This is the basic structure of the model. Its runtime system can divide input, schedule tasks, coordinate communication, and recover unfinished work after machine failures.
A good exposes enough independent work to keep processors busy. A poor creates dependencies that cause workers to wait, exchange excessive data, or repeatedly coordinate.
Takeaway: Divide both work and data deliberately. Independence creates potential speedup, while dependencies and communication determine how much of that potential is realized.
and Communication
means that several activities are in progress during the same period; parallelism means that several activities execute simultaneously on separate processing resources. A single-core server can handle many network requests concurrently by switching among them, while a multicore server can process several requests in parallel. A distributed service may use both properties at once.
Overlapping operations can create nondeterministic outcomes when they access shared state. Important hazards include:
A , in which the result depends on an uncontrolled ordering.
A deadlock, in which processes wait forever for resources held by one another.
Starvation, in which a process repeatedly fails to obtain a needed resource or service.
A lost update, in which one writer overwrites another writer's change after both read the same old value.
Communication between distributed processes usually occurs through messages rather than shared memory. Messages can carry requests, replies, data records, status updates, or coordination signals. The main communication concerns are , bandwidth, reliability, serialization, and backpressure.
A timeout distinguishes a response that has not arrived soon enough from a response that has arrived, but it does not reliably distinguish a failed process from a slow process. Network congestion, processor delay, and message delay can all produce the same outward symptom.
Takeaway: Correct concurrent designs specify which operations may overlap and which must be ordered, protected, retried, or made mutually exclusive.
Coordination, Ordering, and Distributed State
Coordination determines how independent components cooperate, while synchronization constrains the order or timing of operations. Locks and mutexes protect critical sections. Semaphores represent a limited number of permits. Barriers hold participants until all required participants arrive. Queues separate producers from consumers and absorb temporary differences in speed. Leases and heartbeats provide time-limited authority or evidence that a component is responsive.
A is useful when local wall-clock readings cannot safely establish a shared order. It assigns increasing values to events so that a causally earlier event is reflected as earlier in the logical ordering. It does not measure elapsed real time.
A replicated database may require every replica to apply updates in the same order. A can support this requirement by electing a leader, replicating an ordered log, and marking entries committed after enough replicas acknowledge them. Separating leader election, log replication, and safety concerns can make such a protocol easier to reason about.
Distributed state includes each process's local state together with messages currently in transit. No process normally observes this entire state instantaneously. This limitation affects failure detection, snapshots, recovery, replication, and conflict resolution.
describe two different goals. Safety prevents incorrect outcomes, such as conflicting values being committed by two leaders. Liveness ensures that useful progress eventually occurs, such as a request receiving a response. A system that refuses to act during uncertainty may protect safety but reduce availability; aggressive timeouts may improve apparent responsiveness while increasing the risk of incorrect decisions.
Takeaway: Coordination establishes shared decisions, synchronization controls ordering, and explicit goals clarify what the system must guarantee.
and Performance
concerns how performance changes as resources or workload increase. The main measures include speedup, efficiency, throughput, , and tail . For a system using workers, speedup is
where is sequential execution time and is execution time with workers. Efficiency is
Efficiency indicates how effectively the available workers are being used.
Strong scaling keeps the total problem size fixed while adding workers. Execution time may decrease at first, but communication and synchronization eventually dominate. Weak scaling increases the problem size in proportion to the number of workers and aims to keep execution time per worker roughly constant.
is limited by serial work, communication, synchronization, load imbalance, and resource contention. Coordination amplification can also reduce gains: adding workers may increase messages, locks, replicas, and failure possibilities faster than it increases useful computation. Placing computation near the data it needs can reduce network traffic.
Takeaway: More workers do not automatically produce proportional speedup. Performance depends on the balance between useful parallel work and the costs of communication, coordination, and contention.
Failure, Timing, and Fault Tolerance
is a defining challenge of distributed systems. One machine may crash while others continue, a network link may break, a disk may become unavailable, or a process may be reachable from some components but not others.
Common failure types include:
Crash failure: a process stops and performs no further work.
Omission failure: a message or operation is not delivered.
Timing failure: a response arrives too late for the application.
Network partition: groups of processes cannot communicate with one another.
Byzantine failure: a component behaves arbitrarily or sends conflicting information.
Fault-tolerant designs may use replication, retries, checksums, durable logs, quorum rules, failover, idempotent operations, and automated recovery. These techniques reduce the consequences of failure but do not eliminate it. Replicas can diverge, retries can duplicate side effects, and a shared dependency can defeat otherwise independent copies.
Variable timing creates a fundamental ambiguity: a failed process and a merely slow process can look identical to an observer that has not received a response. Protocols therefore depend on explicit assumptions, such as bounded delays, eventual message delivery, or eventual leader stability. Under fully asynchronous message passing, the FLP impossibility result shows that no deterministic can guarantee termination if even one process may crash.
Takeaway: Failure handling must state its assumptions and account for uncertainty. Timeouts, retries, and redundancy are useful tools, but none provides perfect knowledge of remote state.
Integrating the Design Principles
Consider a distributed image-processing service. A coordinator divides a large image collection into partitions, and workers process those partitions concurrently. Messages carry tasks, results, acknowledgments, and retry requests. A queue provides backpressure when workers process data more slowly than producers supply it.
A scheduler monitors heartbeats and reassigns unfinished tasks when a worker becomes unavailable. Results are stored redundantly so that one machine's failure does not destroy them. A replicated metadata service tracks task ownership, while a orders metadata updates. Metrics measure throughput, average , tail , and failed attempts.
Each feature addresses a distinct concern:
supplies independent work.
Communication connects the components.
Synchronization prevents conflicting decisions.
Replication improves durability.
Coordination establishes shared metadata decisions.
Failure handling prevents a local fault from becoming total system failure.
Performance measurement reveals whether added resources produce useful improvement.
The design process can be summarized by six questions:
What work or data can be decomposed?
What operations must be ordered or protected?
What information must be communicated?
How will state be replicated, logged, snapshotted, and recovered?
What happens when a component fails or becomes slow?
How will communication, synchronization, contention, and load imbalance affect scaling?
Final takeaway: Distribution creates both capacity and uncertainty. Strong designs make independence explicit, communicate only what is necessary, define ordering and state rules, and treat delay and failure as normal operating conditions rather than exceptional surprises.