I have no idea where this limit came from. I worked at WhatsApp[1], and while we did split nodes into separate clusters, I think our big cluster had around 2000 nodes when I was working on it.
Everything was pretty ok, except for pg2, which needed a few tweaks (the new pg module in Erlang 23 I believe comes from work at WhatsApp).
The big issue with pg2 on large clusters, is locking of the groups when lots of processes are trying to join simultaneously. global:set_lock is very slow when there's a lot of contention because when multiple nodes send out lock requests simultaneously and some nodes receive a request from A before B and some receive B before A, both A and B will release and retry later, you only get progress when there's a full lock; applying the Boss node algorithm from global:set_lock_known makes progress much faster (assuming the dist mesh is or becomes stable). The new pg I believe doesn't take these locks anymore.
The other problem with pg2 is a broadcast on node/process death that's for backwards compatibility with something like Erlang R13 [2]. These messages are ignored when received, but in a large cluster that experiences a large network event, the amount of sends can be enormous, which causes its own problems.
Other than those issues, a large number of nodes was never a problem. I would recommend building with fewer, larger nodes over a large number of smaller nodes though; BEAM scales pretty well with lots of cores and lots of ram, so it's nicer to run 10 twenty core nodes instead of 100 dual core nodes.
[1] I no longer work for WhatsApp or Facebook. My opinions are my own, and don't represent either company. Etc.
[2] https://github.com/erlang/otp/blob/5f1ef352f971b2efad3ceb403...