You avoided the monotonic key. You hashed tenant_id, as the documentation recommends for a multi-tenant system, and that choice is right about a lot: every tenant's rows live together, so a per-tenant query hits exactly one shard, and tenants spread evenly across the cluster.
Then the marketplace signs one large merchant, and that merchant is 40% of all writes.
What the shards see
The course's lab runs 100,000 writes from 149 tenants across 8 shards partitioned by hash(tenant_id):
| Shard | 0 | 1 | 2 | 3 | 4 | 5 | 6 | 7 |
|---|---|---|---|---|---|---|---|---|
| Share of writes | 2.0% | 11.2% | 49.0% | 7.2% | 7.6% | 7.4% | 11.6% | 4.0% |
| Load vs. average | 0.16× | 0.90× | 3.92× | 0.58× | 0.61× | 0.59× | 0.93× | 0.32× |
Shard 2 carries almost four times the average load while shard 0 sits nearly idle. Capacity is set by the hottest shard, so eight machines deliver roughly the throughput of two.
Hashing distributes keys uniformly. It cannot distribute load uniformly, because load per key is not uniform. The skew is in the traffic, and no better hash function will remove it.
Three repairs, not interchangeable
Salting writes the large tenant's rows under tenant#0 to tenant#15, chosen at random for each write. Sixteen keys spread the load over up to sixteen shards, which here means all eight. The cost lands on reads: every read of that tenant becomes sixteen reads.
A composite key spreads on something the query already knows. Partitioning the large tenant by tenant_id#order_seq evens the load out to 1.01× in the lab, and every query for that tenant now fans out to all eight shards. It is the better choice when queries can supply the extra part of the key.
An explicit placement override gives the large tenant its own shard and leaves the other 148 alone. That shard still carries 40% of all writes by itself, so it needs a bigger machine or a further split. It is the only repair that needs no change to code on the read path.
Think Like an Engineer
The large tenant is salted across sixteen keys, and a report needs that tenant's orders for one day. How many shards does it ask, and what sets its latency? All eight, and the client merges the results, so the report is as slow as the slowest shard's reply. Salting bought an even write load with read fan-out.
You have probably met this problem under other names: "one customer's reports time out and nobody else's do", or "we had to give the enterprise account its own instance". The second one is a hot partition solved by hand.


