Replication and Sharding Together

Sharding divides Shopend’s data across database servers. We can also replicate each shard. Each shard then has its own primary. That primary accepts writes for the shard and passes its changes to the shard’s replicas.

Two routing decisions

For each database operation, the router makes two choices. It looks up the shop in the directory to find its shard. Then it chooses a server within that shard: the primary or one of the shard’s replicas.

The replication choices we made earlier still apply within each shard. For example, after a shopper changes something in Shopend, we send that shopper’s reads to the primary for a few seconds. That is how we provide the read-your-writes guarantee. It assumes the replicas catch up within those few seconds.

Shopend’s database

Here is the whole system. The API load balancer sends a storefront’s request to an API instance. That instance sends its database operations through the shard router:

graph TB
    S[Storefronts] --> LB[API Load Balancer]
    LB --> API1[Shopend API Instance 1]
    LB --> API2[Shopend API Instance 2]
    API1 --> R[Shard Router]
    API2 --> R
    R <-->|Shop-to-shard lookup| D[(Directory)]

    subgraph A[Shard A]
        PA[(Primary)]
        RA[(Replica)]
        PA -.->|Replicate changes| RA
    end
    subgraph B[Shard B]
        PB[(Primary)]
        RB[(Replica)]
        PB -.->|Replicate changes| RB
    end

    R -->|Writes and reads needing the primary| PA
    R -->|Eligible reads| RA
    R -->|Writes and reads needing the primary| PB
    R -->|Eligible reads| RB

Sharding spreads the data across the shards, so each shard’s primary accepts only the writes for its own shops. Together, the primaries can accept more writes than one primary could. Within each shard, the replicas can serve reads, which improves read performance. The replicas also hold copies of the shard’s data, which gives us redundancy.