> hash partitions shouldn't have hot spots unless you hash on a bad column
In the nice world of all data having an even load, sure...
But in the real world, you can easily get a handful of users in your "Users" table sending millions of requests per second, and in that case, you would really like to re-partition so that they don't happen to all end up on the same partition.
Implementation can be as simple as allowing splittable partitions, so that whenever load/size gets too high, a partition can be split in half, and half the records moved to a new host. The partition map is only a few kilobytes, so is globally shared/updated. There are no concurrency issues, because during the splitting process, either the old or new partition is responsible for each record, and both the old and new partition hosts can respond to a read or write query, either themselves or by forwarding it to the other host.
Whenever two neighbouring partitions both see low load/space usage, join them with the same method in reverse. By joining only neighbouring partitions, you can't suffer terrible fragmentation and blowing up the size of the global partition map.