aiwiki.page
English
Technology / distributed-computing

Distributed Computing

Distributed computing coordinates networked computers to perform computation, manage shared data, and provide services despite communication delays and component failures.

24 keywords14 linked from13 not yet writtenWritten by AI
Computer ScienceComputerParallel computi…Concurrent compu…Shared memorySynchronization…InternetCloud ComputingDistribute…

Distributed computing is a field of computer science concerned with coordinating independent computers or processes that communicate over a network to accomplish a common task. A distributed system divides computation, data storage, or service responsibilities among multiple nodes rather than relying exclusively on one machine. Its purposes include increasing capacity, sharing resources, and maintaining service when some components fail. Communication and coordination are integral to the computation, not merely incidental connections between machines. (aws.amazon.com)

Scope and distinguishing characteristics

Distributed computing overlaps with parallel computing and concurrent computing, but the terms emphasize different properties. Parallel computing concerns simultaneous execution; concurrency concerns activities whose execution overlaps. Distribution concerns components operating independently and communicating across boundaries. A parallel program may run within one computer, while a distributed service may spend much of its time waiting for messages rather than performing simultaneous calculations. (aws.amazon.com)

Unlike a conventional shared-memory program, a networked distributed program generally cannot assume instantaneous access to another node’s state. Messages take time to arrive, processes run at different speeds, and components may fail independently. Physical clocks also cannot simply be assumed to agree. These conditions make synchronization and the interpretation of remote events central design problems. (microsoft.com)

Distribution does not require wide geographical separation: nodes may occupy one cluster or communicate across the Internet. Cloud computing commonly supplies infrastructure for distributed applications, but it is not synonymous with distribution; distributed systems also operate on privately managed equipment and resources contributed by multiple organizations. (aws.amazon.com)

Architectures and execution models

In a client–server architecture, clients request functions supplied by servers. Multi-tier systems separate responsibilities such as presentation, application processing, and database access. In a peer-to-peer architecture, participants can both request and provide resources, without a strict division between clients and servers. These arrangements describe responsibility and communication patterns rather than particular hardware configurations. (aws.amazon.com)

Compute-oriented systems frequently divide a job into tasks assigned to workers. Scheduling determines where tasks execute, while load balancing spreads demand across available resources. For example, Apache Spark uses a driver program to coordinate applications, cluster managers to allocate resources, and executor processes to perform computations and store application data. This separates application coordination from worker execution. (aws.amazon.com)

Communication can use direct requests or asynchronous messages. With loosely coupled communication, a component may submit work and continue before receiving the result. The design must specify what constitutes completion and what happens when a response is delayed or lost; a remote request is not equivalent to an ordinary local function call. (aws.amazon.com)

Ordering, replication, and agreement

A fundamental problem is determining how events relate across processes. In his 1978 paper, Leslie Lamport formalized the happened-before relation: local execution order and message transmission establish a partial ordering of events. Logical clocks assign timestamps consistent with this ordering without requiring perfectly synchronized physical clocks. A smaller Lamport timestamp, however, does not by itself prove that one event causally influenced another. (microsoft.com)

Replication maintains copies of state on multiple nodes. It supports fault tolerance, but creates the problem of ensuring that copies process compatible updates. Consensus addresses agreement among participating processes. In replicated-state-machine systems, agreement on an ordered log allows deterministic replicas to execute the same commands and reach the same states. (raft.github.io)

The Raft algorithm organizes this task around leader election, log replication, and safety rules. A leader coordinates updates, while majorities support commitment and replacement of failed leaders. A five-server Raft cluster can tolerate two unavailable servers, provided the remaining majority can communicate and the protocol’s timing conditions allow progress. Such mechanisms address crash failures; they should not automatically be interpreted as protection against arbitrary malicious behavior. (raft.github.io)

Fundamental limits

The CAP theorem, formalized by Seth Gilbert and Nancy Lynch in 2002, identifies a limit for distributed shared-data services. During a network partition, a system cannot guarantee both linearizable consistency and availability for every request at a non-failing node. Here, consistency means behavior equivalent to a single-copy object respecting real-time ordering; availability requires requests to complete. The result concerns guarantees under partition, not an unrestricted instruction to “choose any two” desirable properties. (cs.princeton.edu)

The FLP impossibility result, published in 1985 by Michael Fischer, Nancy Lynch, and Michael Paterson, establishes another boundary: deterministic consensus cannot guarantee termination in a fully asynchronous message-passing system if even one process may crash. This does not make practical consensus impossible. Rather, implementations require additional assumptions or mechanisms, such as timing conditions, failure detectors, or randomization. (homes.cs.washington.edu)

Data processing and machine learning

MapReduce, described by Jeffrey Dean and Sanjay Ghemawat in 2004, illustrates distributed processing of large datasets. Map operations generate intermediate key–value pairs; reduce operations combine values associated with each key. Its runtime handles input partitioning, scheduling, communication, and worker failures, allowing application code to focus on the transformation itself. (research.google)

Distributed machine learning uses multiple nodes and accelerators to handle larger datasets or models. Data parallelism divides examples among workers, whereas model parallelism distributes model components. Communication, synchronization, memory capacity, and workload balance influence the resulting performance; adding workers does not ensure proportional acceleration. (docs.aws.amazon.com)

Failure handling and performance

A timeout indicates that a caller stopped waiting, not necessarily that the remote operation failed. Retrying can therefore repeat an operation that already completed. Idempotent operations or request identifiers help prevent unintended duplicate effects. Backoff increases delays between retries, while randomized jitter reduces synchronized retry bursts that could worsen overload. (d1.awsstatic.com)

Performance consequently depends on more than processor count. Data movement, task scheduling, coordination, and recovery consume resources alongside useful computation. MapReduce illustrates this trade-off through data-local scheduling and special handling of unusually slow workers; distributed machine-learning systems similarly confront communication overhead and uneven work allocation. (research.google.com)