Practical systems deal with this by not caring strongly about overflow (caches), by running at low enough utilizations that overflow is very unlikely given their item counts (e.g. Dynamo), by using explicit partitioning rather than consistent hashing (e.g. DynamoDB), by being able to take advantage of multi-tenancy to drive up per-physical-node utilization even in the case of low per-logical-node utilization, or by using some additional algorithmic sophistication (e.g. Chen et al https://arxiv.org/pdf/1908.08762).
In practice, this kind of overflow is a big deal for systems that deal with relatively small numbers of large objects, and are not as big a deal for systems that deal with large numbers of small objects. Try out the numbers in the page's "Handy Calculator" to see how that plays out.
It's also worth mentioning that this isn't unique to consistent hashing, but is a problem with random load balancing more generally. "Pick a random server and send traffic to it" is an OK load balancing strategy when requests are small and servers are large, but a terrible one when requests become relatively large or expensive. In the general load balancing/placement problem this is easier than the storage case, because you don't need to find requests again after dispatching them. That makes simple algorithms like best-of-2 and best-of-k applicable.