Replication
Replication, also known as data redundancy, is the process of creating and maintaining identical copies of data across multiple servers or nodes.
In a monolithic architecture, data is typically stored in a single database instance, and joins are typically used to query relational data across multiple tables.
But in a distributed software, those joins turn into network calls. This increases latency, decreasing performance, and also having implications for the resilience of the system.
Data replication is one possible solution to these problems. Simply, data is replicated to wherever it is needed, so that discrete services are less reliant upon each other. Services therefore work faster by needing to use only local data.
The more that an application becomes distributed, for example in service-oriented or microservices architectures, the more necessary it becomes to replicate data across multiple servers.
Besides the performance benefits of reducing network calls, data replication also improves the fault tolerance and availability of distributed software, due to reduced dependencies between services. Replicas also spread read load across nodes, so a hot primary is no longer the single bottleneck for queries, raising read throughput without changing the write path.
In cloud platforms, replicas are commonly spread across availability zones to survive data-center-level failures.
Replication, combined with sharding, is how distributed databases provide both high horizontal scalability and fault tolerance. Sharding splits data across multiple servers, while replication duplicates data across multiple servers. Combining the two allows a cluster to grow in capacity and to survive node losses.
Data need not be replicated in the form of its source-of-truth. A replica can hold a transformed or denormalized projection rather than a verbatim copy. Materialized views are the common case, storing pre-processed data that’s optimized for specific use cases, such as reporting. Such data replication strategies are commonly use for performance optimization.
Besides improving scalability and reliability, another use case for data replication is to support migrations, re-platforming, and other architectural evolution. For example, in the transition from monolithic services to more microservices, data replication can help to reduce the need for scheduled downtime while data migrations are run. A new service can be run against a replicated copy of the monolith’s tables while the cut-over is validated, so the legacy system keeps serving traffic. Once the new service is trusted, it takes ownership of its data and the legacy copy is retired, moving toward one source-of-truth per service.
Replication strategies
There are three broad replication strategies. They differ in where writes are accepted and how the resulting copies are kept in sync.
- Single-leader replication
- Multi-leader replication
- Leaderless replication
Single-leader replication
Only one node, the leader, accepts writes. One or more replicas apply the leader’s changes asynchronously, and can serve reads alongside it. Clients may read from the leader or from any replica, but every write must go through the leader.
This is the simplest topology to reason about, and gives strong write consistency. The trade-off is the leader is a single point of failure, unless external failover logic promotes a replica in its place, which itself requires reconfiguring the cluster.
Because replication is asynchronous, replicas typically lag behind the leader, so a read from a replica can return stale data, and a leader under heavy write load can become a bottleneck.
Single-leader replication suits systems that need read scalability without giving up write consistency or simplicity. OLTP systems, dashboards, and APIs with moderate write throughput are common cases.
Multi-leader replication
More than one node accepts writes, and each leader replicates its own changes to the others — write locally, sync globally.
Because two leaders can accept conflicting writes to the same record, multi-leader setups need an explicit conflict-resolution policy, such as last write wins, version vectors, or CRDTs (conflict-free replicated data types).
This buys high availability across regions and lets writes commit with local latency, at the cost of harder conflict resolution, increased cross-node replication latency, and a data state that’s harder to reason about.
Multi-leader replication suits systems that must keep accepting writes across regions, or during partial outages. Collaborative tools, messaging platforms, and edge-aware databases are common cases.
Leaderless replication
Used by systems such as Cassandra and DynamoDB, leaderless replication has clients write directly to multiple nodes rather than through any single leader. Consistency is instead managed through quorum reads and writes. A write is acknowledged once W nodes confirm it, and a read queries R nodes, with W + R chosen greater than the replica count so that every read is guaranteed to see at least one up-to-date copy.
Leaderless systems trade away a designated leader for very high availability, at the cost of eventual consistency, unless the quorum sizes are tuned carefully, and consistency guarantees that are generally harder to reason about than a single-leader system’s.
Trade-offs
Trade-offs are made in terms of data consistency, and the increased system complexity required to manage the data replication and synchronization processes. Storage costs also increase, since every replica is another full copy of the data, and the bandwidth needed to keep copies in step grows with the write rate.
The consistency trade-off is the heart of the CAP theorem. Once data lives on more than one node, a network partition can force a choice between refusing writes to keep copies in step, or accepting them at the cost of temporary divergence between replicas. Most replicated stores lean toward the latter, letting copies converge over time rather than blocking writes while a partition heals.
Data synchronization can be achieved using either batch processes or event streams. Batch processes periodically dump and reload snapshots of the source data, which is simple to reason about but leaves replicas stale between runs. Event streams propagate each change as it happens, so replicas track the source with far less lag, at the cost of ordering, delivery, and replay machinery. Some database systems also provide built-in support for data replication across multiple nodes, typically by shipping their write-ahead log or replaying statements row by row. Built-in replication is transparent to applications but usually ties the replicas to the same database engine as the source.
Leaderless and eventually-consistent stores often synchronize replicas with an anti-entropy gossip protocol, which repairs divergence through randomized peer-to-peer exchanges.