Dividing the Data
Earlier, we used replication to spread Shopend’s reads across database servers. Caching serves some popular reads, and replicas serve other reads that can tolerate replication lag. But all application writes still go to the primary.
Our assumed workload was 10,000 reads and 1,000 writes per second. Moving eligible reads away from the primary leaves it more capacity for those writes. Suppose the write load keeps growing until the primary cannot keep up. Adding another replica does not divide that work: the primary still accepts every write. Moreover, as the write load increases, the work on replicas increases too, because every replica applies every change made to the primary to keep its copy up to date.
An alternative to replication is to divide the data, so that different servers handle the reads and writes for different parts of the data.
Sharding
Dividing a dataset across database servers is called sharding. Each part is a shard. A shard stores part of the application data and handles reads and writes for that part.
Suppose we have nine records and three database servers. With replication, all three servers hold records 1 through 9. With sharding, we could put records 1 through 3 on one server, 4 through 6 on another, and 7 through 9 on the third:
graph LR
subgraph Replication
R1[(Replica A<br/>Records 1–9)]
R2[(Replica B<br/>Records 1–9)]
R3[(Replica C<br/>Records 1–9)]
end
subgraph Sharding
S1[(Shard A<br/>Records 1–3)]
S2[(Shard B<br/>Records 4–6)]
S3[(Shard C<br/>Records 7–9)]
end
Replication ~~~ Sharding
A change to record 2 goes to shard A. A change to record 8 goes to shard C. Since shard B does not store record 2 or 8, it does not need to apply these changes.
Sharding and replication can be used together. Each shard can have replicas that hold copies of its part of the data. Those replicas can serve reads that tolerate replication lag. They can also serve as failover options: one of them can take over if that shard’s primary fails.
Shard router
With sharding, we cannot send an operation to whichever server happens to be available. A read (or write) for record 2 must reach shard A, because shards B and C do not hold it.
A shard router directs each database operation to the shard(s) holding the data it needs:
graph TB
A[Shopend API Instances] --> R[Shard Router]
R --> S1[(Shard A)]
R --> S2[(Shard B)]
R --> S3[(Shard C)]
The router can be a component that comes with the database, or it can be routing logic in the application’s data access layer. The diagram shows what the router does. It does not mean the router has to run on its own server.
A shard router is similar to an API gateway, which picks a service based on the operation that service handles. The difference is that a shard router picks shards based on the data they hold.