A distributed system is a computing architecture in which multiple independent machines coordinate over a network and appear to users as one service. Examples include streaming platforms and e-commerce backends. Every Kubernetes cluster also fits this definition. The moment you spread work across machines, though, you inherit network failures plus coordination and cross-service debugging overhead that a single server never has.
In brief:
- A distributed system coordinates multiple machines that share no memory and communicate over unreliable networks, yet appear to users as one service.
- The CAP theorem forces a choice during network partitions: refuse requests for consistency and partition tolerance (CP) or serve possibly stale data for availability and partition tolerance (AP).
- Common patterns include cluster computing, cloud platforms, content delivery networks (CDNs), peer-to-peer networks, and distributed databases, each tuned for different trade-offs.
- A self-hosted Strapi 5 deployment can use external persistent data and media in a custom multi-replica architecture, with auto-generated Representational State Transfer (REST) application programming interfaces (APIs) and configured webhooks for downstream cache and CDN invalidation.
Distributed vs. centralized systems: when to choose each
Distribution earns its overhead when you need fault isolation, horizontal scale, or geographic reach that one machine cannot deliver. If your application fits on a single server and strong consistency is non-negotiable, staying centralized saves real complexity.
In Stack Overflow's 2016 architecture write-up, the team reported 209,420,973 Hypertext Transfer Protocol (HTTP) requests in a single day while running its entire question-and-answer network off a single application pool on a single server, a useful corrective to the assumption that scale always requires distribution.
| Aspect | Centralized Systems | Distributed Systems |
|---|---|---|
| Scaling | Vertical (bigger hardware) | Horizontal (more machines) |
| Failure impact | Single point of failure | Isolated failures, graceful degradation |
| Consistency | Strong guarantees, Atomicity, Consistency, Isolation, and Durability (ACID) transactions | Often eventual consistency |
| Debugging | One stack trace tells the story | Tracing across service boundaries |
| Deployment | Single process | Independent service deployments |
| Network dependency | None between components | Network failures affect functionality |
| Cost | Lower operational complexity | Higher tooling and operations cost |
Distribution can deliver four benefits:
- Horizontal scalability: add or remove nodes as demand changes, matching provisioned capacity to current load.
- Fault tolerance: redundancy can be structural, allowing the service to withstand some node failures before users notice.
- Geographic performance: placing compute and data near users can cut round-trip latency and spread risk across regions.
- Elastic cost: usage-based infrastructure can align spending more closely with actual demand than permanent peak-capacity provisioning.
These benefits come at the cost of coordination complexity, which the rest of this article unpacks.
Distributed systems vs. microservices
Microservices architectures are one type of distributed system. A replicated monolith behind a load balancer is distributed too and faces the same consensus, consistency, and partition problems without any service decomposition. Microservices orchestration adds an organizational dimension on top: independently deployable services with firm boundaries, owned by separate teams.
Martin Fowler's microservices advice is blunt: don't even consider microservices unless you have a system that's too complex to manage as a monolith. The failure mode to watch for is the distributed monolith, a system packaged as services that still behaves like one tightly coupled application stretched across the network, with lockstep releases and long synchronous call chains. You get microservices' operational cost without their independence.
How distributed systems work: nodes, network, and coordination
Five traits distinguish distributed systems from single-server designs:
| Trait | Meaning |
|---|---|
| Resource sharing | Compute and data on one node support workloads on another. |
| Concurrency | Multiple users access resources simultaneously, sometimes coordinated with locks and queues. |
| Scalability | Capacity grows by adding nodes. A centralized system scales by upgrading one machine. |
| Transparency | Clients see one service, not the individual nodes or their partial failures. |
| Fault tolerance | The service can keep running through some hardware failures or dropped packets when sufficient redundancy exists. |
In practice, distribution can reduce the need for emergency hardware upgrades during traffic spikes and keep services reachable through some partial outages. The cost is complexity in every layer, and the CAP theorem (below) requires an explicit trade-off during network failures.
Physical servers, database instances, and an API server can all be nodes that run workloads or store data. Fault-tolerant designs can treat them as replaceable.
Network is the communication fabric. Bandwidth limits and latency spikes live here. So does packet loss.
Coordination handles consensus and scheduling. Health checks also help turn independent machines into what users perceive as one service.
The network layer is where many designs get into trouble. The Fallacies of Distributed Computing, attributed to Peter Deutsch with an eighth added later by James Gosling, begin with "the network is reliable," and that assumption still causes production failures. When you design replication or sharding (partitioning data across nodes), you are accounting for network behavior as well as application logic.
Coordination is where consensus protocols operate. The Raft consensus protocol breaks consensus into leader election and log replication while enforcing safety guarantees: a follower that stops receiving heartbeats within its election timeout starts an election, and the first candidate to win votes from a majority of servers becomes leader for that term.
Randomized election timeouts prevent split votes. You rarely implement this yourself. etcd coordination primitives include distributed locks and elections, and Apache ZooKeeper's recipes implement leader election with sequential ephemeral znodes (ZooKeeper data nodes), where the process holding the smallest sequence number leads.
Types of distributed systems and when each fits
Cluster and grid computing
Cluster computing runs machines on the same low-latency network as one logical supercomputer and slices tasks into parallel jobs for compute-intensive workloads. Google's Borg cluster manager follows this pattern and preceded Kubernetes. Kubernetes is now a widely used way for teams to run clustered workloads: the Cloud Native Computing Foundation's 2025 CNCF Annual Survey found 82% of container users run Kubernetes in production, up from 66% in 2023. Teams looking at Kubernetes for CMS workloads specifically can see an example in Strapi's guide to deploying on Kubernetes.
Grid computing can loosen the coupling: geographically scattered resources, often owned by separate organizations, donate capacity to problems no single cluster could handle. The European Organization for Nuclear Research (CERN) operates the Large Hadron Collider (LHC), whose Worldwide LHC Computing Grid spans over 170 sites in 42 countries with about 1.4 million computer cores, and Folding@home protein-folding work has been distributed across volunteer machines for more than 25 years.
Cloud computing
Infrastructure as a service (IaaS) provides virtual machines, and platform as a service (PaaS) provides managed runtimes. Software as a service (SaaS) provides complete applications. Cloud platforms deliver this elasticity as a utility. Amazon Web Services (AWS), Microsoft Azure, and Google Cloud Platform (GCP) can reduce server-management and capacity-planning work, while capabilities such as global failover depend on the services and architecture selected. Gartner's November 2024 forecast put worldwide public cloud spending at $723.4 billion for 2025, up from $595.7 billion in 2024. Strapi Cloud is one example of a PaaS option purpose-built for CMS workloads, handling infrastructure and deployment options so you can focus on content modeling.
Content delivery networks (CDNs)
Edge servers cache static assets, API responses, or entire pages close to users to cut round-trip latency and absorb traffic spikes. One caveat matters for API-driven applications: CloudFront API caching requires activation for API Gateway origins, while Google Cloud CDN excludes JavaScript Object Notation (JSON) by default in CACHE_ALL_STATIC mode. Caching API responses requires explicit configuration and an invalidation strategy. This caching pattern can improve headless CMS performance, and a Strapi REST cache plugin is available for more granular control.
Peer-to-peer systems
Peer-to-peer systems distribute data directly among participating nodes. In BitTorrent file sharing, participants downloading the same file upload pieces to one another.
Distributed databases
Distributed databases partition and replicate data across nodes for scale, availability, or both. The replication model determines behavior during failures:
- CockroachDB replication works synchronously via Raft. A write commits only after a quorum of replicas confirms it, so writes to a range that cannot reach a quorum cannot commit.
- Cassandra replication is leaderless and multi-primary: any node accepts writes, timestamped last-write-wins resolves conflicts, and hash-based Merkle trees efficiently compare replica contents for background repair.
- DynamoDB read consistency defaults to eventually consistent reads and lets you request strongly consistent reads per operation. Eventually consistent reads cost half as much in read units, so the consistency choice shows up directly on your bill.
These databases also differ in how they partition data across nodes.
DynamoDB partition placement uses hashed partition keys. CockroachDB range partitioning splits the key space into contiguous ranges. Geographic partitioning places data near its origin for latency and data-residency requirements.
The CAP theorem: the core distributed system trade-off
The CAP theorem, formalized by Gilbert and Lynch in 2002, tells you that during a network partition a distributed system can guarantee at most two of three properties:
- Consistency: every read returns the latest written value (linearizability).
- Availability: every request received by a non-failing node must result in a response.
- Partition tolerance: the system keeps operating when the network loses messages between nodes.
CAP's "C" means single-copy consistency, while ACID's "C" means preserving database rules. Brewer's 2012 CAP refinement states it directly: CAP refers only to single-copy consistency, a strict subset of ACID consistency, while ACID's C means a transaction preserves database rules such as unique keys.
Since multi-region partitions are inevitable, the practical choice is between CP and AP behavior. The deciding question is what your users tolerate better: stale reads or failed writes.
CP behavior during partitions
During partitions, CockroachDB preserves consistency by refusing requests or returning errors, preventing stale reads. ZooKeeper also prioritizes consistency during partitions. This CP behavior fits leader election, configuration management, and anything where correctness is non-negotiable.
AP behavior during partitions
Cassandra stays available and converges eventually. Abadi's PACELC classification, meaning that a system chooses availability or consistency during a partition and latency or consistency otherwise, puts Dynamo-style systems in the PA/EL category. This AP behavior fits always-on global applications where latency matters more than immediate consistency.
Brewer's 2012 paper also warns you against treating the labels as rigid: the choice between C and A can occur many times within the same system at very fine granularity.
Kleppmann's critique goes further, cautioning that labeling whole databases CP or AP obscures more than it explains. The labels describe behavior during partitions and do not categorize products. Daniel Abadi's PACELC extension adds that even without a partition you trade latency against consistency on every operation.
The BASE model (Basically Available, Soft State, Eventually Consistent) is the application-level expression of the AP choice. As Dan Pritchett put it in ACM Queue: where ACID is pessimistic and forces consistency at the end of every operation, BASE is optimistic and accepts that the database consistency will be in a state of flux. Pritchett also notes that in partitioned databases, trading some consistency for availability can lead to dramatic improvements in scalability.
Common distributed systems challenges in production
Network partitions are routine, not rare
Bailis and Kingsbury's network reliability study in Association for Computing Machinery (ACM) Queue documents that one company running 100-200 nodes on a major hosting provider saw five distinct partition periods in a 90-day window.
The same survey records a MongoDB replica set on Amazon Elastic Compute Cloud (EC2) where a partition separated the primary from its secondaries. When the old primary rejoined two hours later, it rolled back everything written on the new primary. Two hours of writes, gone. A network call behaves differently from a function call, and designs that treat them alike eventually pay for it.
Split-brain and consensus
Split-brain data corruption occurs when a partition lets both sides of a cluster elect a leader and accept writes. Google's site reliability engineering (SRE) book recommends formally proven consensus systems for distributed locking. The recommendation also applies to leader election and critical shared state, because informal approaches to solving this problem can lead to outages, and more insidiously, to subtle and hard-to-fix data consistency problems.
Coordination services create their own dependency risk. Google's Chubby lock service became so reliable that dependent teams assumed it would never fail, so Google introduced planned Chubby outages to stop services from over-relying on availability that was never guaranteed.
Debugging across service boundaries
In a monolith, a stack trace often tells the whole story. In a distributed system, you reconstruct what happened by tracing requests across services, and correlation ID propagation is the foundation:
const correlationId = req.headers['x-correlation-id'];
logger.info({ correlationId, service: 'checkout' }, 'Checkout started');OpenTelemetry standardizes this propagation via World Wide Web Consortium (W3C) Trace Context headers, and its tracing specification is stable with long-term support. Jaeger v2 is now built directly on the OpenTelemetry Collector framework, and its legacy client SDKs are retired in favor of OpenTelemetry software development kits (SDKs). Traces are usually sampled, and OpenTelemetry sampling notes that high-volume systems commonly sample 1% of requests or fewer, so you rarely have every request on record. Instrumenting early matters; retrofitting tracing after an incident is the hard way to learn this, a pattern that also shows up in Strapi and Next.js performance.
Operational cost
Distributed systems demand investment in monitoring tools and operational expertise, and the observability stack meant to reduce operational burden brings its own scaling and cost requirements. Quorum-based operations require communication with multiple nodes, so coordination can add network latency to the request path.
Patterns that make distribution survivable
Give every network call a timeout, then retry with exponential backoff and jitter; without jitter, clients that failed together retry together and hammer a recovering service into a thundering herd.
You can make the operations behind those retries idempotent, usually by having the caller send an idempotency key that the server records before doing the work, so a retried request cannot charge a card twice or create two orders.
You should assume your message broker delivers at least once, which puts deduplication on the consumer through storing processed message IDs and dropping repeats.
For business transactions that span several services, the saga pattern breaks the work into local transactions with a compensating action for each step. A failure at step four triggers refunds and cancellations through compensation; global lock rollback belongs to a different transaction model.
Where Strapi fits in a distributed architecture
Strapi is an open-source headless CMS built on Node.js and serves as the content service in a distributed architecture. Strapi 5 auto-generates REST endpoints for every Content-Type you define, so the API surface tracks your content model. Those endpoints are private by default until you allow the relevant actions for the Public role or issue an API token. GraphQL requires installing the official GraphQL plugin; the GraphQL vs REST comparison for Strapi covers when each fits.
For self-hosted deployments, the officially documented pieces map onto the layers above:
- Strapi exposes a
/_healthendpoint that returns HTTP 204, suitable for load balancer health probes. - Persistent state lives in a supported database: PostgreSQL 14+, MySQL 8+, MariaDB 10.3+, or SQLite, per the deployment requirements. MongoDB is not supported, and Strapi's serverless guidance says it is not suited to serverless environments because of cold-start performance.
- Media files can move to external storage after you install and configure an upload provider. The Amazon S3 provider is officially maintained and works with S3-compatible services like Cloudflare R2 and MinIO; Azure Blob and Google Cloud Storage providers are community-maintained.
These pieces can support a custom multi-replica deployment. Native Strapi orchestration is not part of this setup. A Strapi Kubernetes deployment provides an example beyond Strapi 5's core deployment docs: it runs three Next.js replicas behind a load balancer with CloudFront in front. It uses Redis publish/subscribe (Pub/Sub), so each Next.js pod invalidates its own application cache when content changes.
Replica management and load balancing are hosting-platform work. The frontend handles framework-cache revalidation, while Redis distributes invalidation messages among frontend replicas. Strapi's webhook system can drive these invalidation flows by notifying downstream services whenever content events fire.
The CDN requires its own purge or caching policy. Configured Strapi webhooks can notify downstream services of content events, but each downstream cache, whether a rendering framework or a CDN, needs its own invalidation path. On Strapi Cloud, API-response caching is opt-in through Cache-Control headers.
Strapi Cloud manages hosting infrastructure and provides managed PostgreSQL, a built-in global CDN, automatic production database backups that do not include assets or secondary environment databases (weekly on Pro, daily on Business), and a 99.9% uptime service-level agreement (SLA) on Business, per Strapi Cloud pricing. For teams evaluating their overall headless architecture, Strapi Cloud removes the infrastructure layer so you can focus on content modeling and API design.
Deciding whether you need a distributed system
You can start with what your users can tolerate during a failure: stale reads or refused writes. That answer picks your side of the CAP trade-off, which in turn narrows your database and coordination choices. Then consider distributing only the layers that need it. A self-hosted Strapi deployment with persistent data and media stored in external services can support a custom multi-replica architecture based on the stateless process pattern, in which no instance keeps local persistent state.





