e.g. we have 4 shards, so user 5 is on shard1. If we go to 6 shards, user5 is now on shard5.
I guess it works with downtime to move the users, or another layer of indirection, where the newly created shards can "point back" to existing shards, but otherwise I don't see it.
My understanding of sharding is that it's best to have a global lookup table user->shard. That's not a huge amount of data, even for millions of users.
Anyone care to educate me as to how they get user%num_shards working in practice?
In practice it's going to be very rare that you'll overwhelm a single shard and adding an extra layer to point to where people actually are is quite simple and fast. You can use a simple cache key (read-through cache of course) that, if it exists means the user is on a specific shard, overriding the default algorithmic pick.
12 shards on 1 machine
6 shards on 2 machines
4 shards on 3 machines
3 shards on 4 machines
2 shards on 6 machines
And on and on...Edit: no wait. If you can split a shard across multiple machines, what's the benefit of having more than 1 shard? Why not have 1 shard split across 1000 machines?
In order to split across multiple dbs, you're looking at creating say 100/1000 dbs in our initial split (when you've got maybe 2-3 machines). And that number then caps the number of machines you can scale to without adding another layer (sharding-shards) or having downtime?
Yes, but you say that like it's a bad thing, it's not. If you're ever forced to reshard, that means you grew beyond what you ever hoped... hurray, awesome, nice problem to have. The reality is however, 99% chance that'll never happen.
Like 1 in a 1000 people ever run into a scaling problem of this magnitude but if you read the blogs you'd come away thinking scaling issues that require sharding are common and everyone needs this stuff, but they aren't, and they don't.
Presumably consistent hashing is also helpful here.