
Distributed system technologies are the protocols, software, data stores, and operational tools that help separate computers communicate and coordinate work. They include remote procedure calls, message brokers, service discovery, databases, coordination services, and observability tools.
The useful question is which problem each technology solves. An application may need a queue for background work without needing a container orchestrator. A replicated database may handle coordination internally without the application implementing a consensus algorithm. This guide explains those choices and the failures they still leave you responsible for.
A distributed application has components that exchange messages across a network. Those components can run under one organization or several, in one location or many. Distribution describes how work and state are arranged; it does not establish who controls the system. Our guide to distributed and decentralized systems explains that separate question.
The main technology families serve different purposes:
There is no required combination of products. Choose the smallest set that meets the workload's capacity, recovery, data, and operating requirements. Adding machines can increase capacity or provide alternate execution paths, but only when the application and its dependencies can use them.
A technology is something you adopt, such as an RPC library, database, or message broker. A technique describes behavior: replication copies state, partitioning divides it, and idempotency makes repeated execution have the same intended effect as one execution.
One product can implement several techniques. A database may partition records, replicate each partition, and use consensus to agree on updates. Conversely, a team can implement a caching strategy using an in-process library or a separate cache service. The appropriate choice depends on how the application reads and changes information.
Keep an explicit record of the guarantee you need. “Use a queue” leaves open whether a job can be lost, delivered twice, or processed out of order. “Retry failed requests” leaves open whether a timeout means the first attempt failed at all.
Transport protocols move data between endpoints. TCP provides a reliable, ordered byte stream while a connection remains usable. UDP sends datagrams without providing those delivery and ordering guarantees itself. QUIC builds a secure transport over UDP with features including reliable streams; HTTP/3 uses QUIC.
These layers should not be treated as interchangeable competitors. HTTP defines application-level behavior, while TCP and QUIC provide transport. Nor is UDP automatically the fastest choice for an application: any reliability, congestion control, or recovery the application needs must still come from somewhere.
Choose supported protocols first, then measure representative payloads, connection behavior, and network conditions. A successful transport exchange does not prove that the remote application committed a database update. For the broader topology and failure-domain decisions, see our distributed network architecture guide.
An API specifies how components request work or exchange information. An HTTP API can expose resources or operations; an RPC framework presents remote service methods to client code. Neither approach removes the network boundary.
gRPC, for example, defines services and methods and commonly uses Protocol Buffers to describe messages. It supports individual request-response calls and several streaming patterns. Client libraries handle much of the communication machinery, while applications still decide what an operation means.
A remote call can finish on the server after the client stops waiting. Retrying it may repeat an effect. For a payment or order operation, use a request identifier and a durable record of the outcome, with concurrency handling that prevents two attempts from applying the effect twice. A key alone does not establish that behavior.
Set a realistic deadline and pass the remaining time budget through dependent calls. The gRPC deadline guidance also distinguishes canceling a call from stopping the work an application has started. Cancellation does not undo changes already committed.
JSON, Protocol Buffers, Avro, and MessagePack are examples of data formats used between components. Evaluate payload size, supported languages, debugging needs, schema rules, and compatibility between deployed versions. Binary encoding alone does not make an API faster or safer.
The Protocol Buffers language guide identifies changes that are safe or unsafe for its binary format. Application meaning still needs separate care: a field that remains readable can acquire an incompatible business meaning. Test old clients against new services and new clients against old services before rolling out a change.
A queue can hold work until a consumer is ready. Publish/subscribe lets several interested consumers receive messages. A retained event log also lets consumers revisit earlier events within the retention policy. These patterns suit background jobs, integration, and workflows that do not require every step to finish before the caller receives a response.
Apache Kafka organizes events into topics and partitions. Ordering is defined within a partition, so choosing a partition key affects which related events share an ordering boundary. Retention and consumer positions allow replay, but application behavior must remain safe when earlier events are processed again.
Separate acceptance by the broker from completion by the consumer. RabbitMQ's acknowledgement documentation explains why publisher confirms and consumer acknowledgements cover different parts of delivery. Neither is a blanket promise that a downstream business operation happened once.
Before adopting messaging, define duplicate handling, ordering needs, retries, retention, and what happens to a message that repeatedly fails. Bound the backlog and monitor consumer lag. A queue absorbs a temporary mismatch in processing rates; it cannot make a permanently undersized consumer keep up.
Service discovery maps a logical name to endpoints as instances start, stop, or move. A small stable deployment may use configured addresses or DNS. More dynamic deployments may use a registry or a platform's built-in discovery.
Kubernetes Services and EndpointSlices provide one example: a Service represents a set of endpoints, and EndpointSlices record those backends. The platform can update them as workloads change. Finding an endpoint still does not prove that a particular request will succeed.
Load balancing selects among suitable destinations. Round robin, connection counts, locality, or application-specific keys can inform that choice. Equal request counts may produce unequal load when requests have different costs. Health checks, connection draining, and retry behavior also affect the result.
Use routing that matches the application's state. Sending a request to another instance is useful only if that instance can access the required session or data. A load balancer does not move that state for you.
Applications may use a relational database, key-value store, distributed database, object store, or shared filesystem. Choose from access patterns and guarantees: transactions, query shape, object size, concurrent updates, acceptable read staleness, and recovery requirements.
A distributed application does not require a distributed database. Several services can use a database on one machine, although that database then becomes a shared capacity and availability dependency. A managed database also has its own internal design, which may include replication that customers do not operate directly.
Separate durability, availability, consistency, and backup. A service can retain data safely while temporarily refusing requests. Replicas can stay available while lagging behind a writer. Replication can also copy an accidental deletion. Backups need their own retention and restore checks.
Technology cannot decide which team owns a customer record or reconcile two different definitions of an order. Those concerns belong to the broader distributed information system, including schemas, integration, and information ownership.
Replication maintains copies of data or state. The update protocol determines when a write is acknowledged and what readers can observe. Synchronous and asynchronous arrangements make different trade-offs between waiting, availability, and potential data loss after a failure. PostgreSQL's replication documentation gives a concrete example of these choices.
Partitioning, also called sharding in many database contexts, divides data or work among components. A partition key may use a tenant, customer, hash, or range. It can spread work beyond one machine while creating hotspots, cross-partition transactions, or expensive rebalancing.
For an order service, partitioning by customer might keep a customer's orders together but leave a large customer dominating one partition. Test the expected distribution of activity, not just an evenly generated benchmark. Adding partitions does not automatically fix a poorly chosen key.
Consensus lets participating nodes agree on a decision or an ordered sequence of state changes under a defined failure model. It supports tasks such as leader election, membership, and replicated configuration. Applications usually consume these capabilities through a database or coordination service.
etcd uses Raft and requires a majority of voting members to agree on updates. A five-voter cluster needs three votes. If a partition leaves groups of three and two, only the group with a reachable majority can continue committing updates, assuming the other operating requirements hold.
Consensus cannot keep every isolated group writable while preserving the same agreement guarantee. It also does not decide application rules, guarantee a deadline, or prevent an authorized client from submitting a bad update. Most teams should use a tested implementation and understand its recovery procedures before relying on it.
A cache holds information that is cheaper to retrieve than to recompute or reload. Options include an in-process cache, a separate service such as Redis or Memcached, and an HTTP cache close to users. Each has a different sharing and failure boundary.
Define how entries expire or are invalidated and which operations must bypass the cache. An old product description may be acceptable for a short period; an old authorization decision may not be. Time-to-live limits how long an entry remains without refresh, but does not guarantee that it matches the source during that interval.
Plan for cache loss and simultaneous misses. If many requests recompute the same value at once, the backend can become overloaded. Coalescing that work, limiting concurrent refreshes, and testing operation without the cache are often more useful than optimizing hit rate alone.
Containers package applications with their user-space dependencies. Orchestrators such as Kubernetes place workloads, maintain a desired number of running instances, and support deployment and networking operations. These capabilities can help when many workloads change independently.
They do not supply the application's consistency or recovery design. Restarting a process does not establish whether its last database write completed. Running several replicas does not make in-memory session data available to all of them. Configuration and persistent storage still need deliberate handling.
A distributed system can run on physical machines or virtual machines without containers. For a small deployment, simpler process management and deployment automation may meet the requirements with less operational work. Adopt orchestration when its scheduling and lifecycle capabilities justify operating the control plane or depending on a managed one.
Logs record events, metrics describe measurements over time, and traces connect spans of work across instrumented components. Together they can help explain slow requests, errors, resource pressure, and waiting between services.
OpenTelemetry supplies instrumentation and mechanisms to generate, collect, and export telemetry. You still need a destination to store and analyze it. Trace context must propagate across the boundaries you want to follow, including asynchronous messages where appropriate.
A trace can show that an inventory call took most of a request's time. It does not prove that every replica contains the same inventory value. Sampling, missing instrumentation, and broken context propagation also limit what is visible. Use explicit application checks for data invariants and recovery behavior.
Choose alerts around user outcomes and actionable failure signals. Decide who responds and which evidence they need. Keep credentials and unnecessary personal data out of telemetry, and control access to the records you retain.
Consider an illustrative order service. This is one possible flow, not a required architecture:
The fourth step illustrates a transactional outbox. It avoids leaving an order committed without a corresponding publication record. A relay can still publish twice if it fails before recording progress, so consumers need duplicate handling. Recovery requires checking the full path, including external side effects.
Start with the reason for distribution and the behavior users need during failure. Then work through the decisions in an order that makes the technology choices testable:
During testing, verify business outcomes as well as process health. After a lost response and retry, there should still be one intended order. After a consumer restart, pending work should resume without silently skipping items. After restoring a backup, the application should be able to use the recovered data.
Cloud computing describes how computing capabilities are provisioned and consumed as services. Distributed system technologies describe how separate components communicate and coordinate. A cloud platform may operate a broker or database for you, while your application still owns message meaning, access configuration, and its use of service guarantees.
You can build a distributed application on owned servers, rented virtual machines, or a mixture of locations. You can also consume a cloud service without designing its internal distributed system. Our cloud and distributed computing comparison explores that distinction.
If one application and database meet your capacity, recovery, and location requirements, that may be the best starting point. A background worker, a standby database, or a second application instance can address a specific constraint without introducing a large collection of independently deployed services.
Those additions can themselves create distributed behavior, so simplicity does not remove the need to test failure. It limits how much coordination the team must understand. Before adopting another technology, name the constraint it addresses, demonstrate the benefit with the real workload, and identify who will maintain it.
Pick one AI, compute, or storage workload and see the difference for yourself. Spin it up in minutes, or let our team map your fastest path to production.