6 ms·
The overlap between K8s and BEAM is a good question. Even amongst experienced BEAM (especially Erlang) programmers, there's a lot of conflicting information. Fr
by dynamite-ready 6y ago
The overlap between K8s and BEAM is a good question. Even amongst experienced BEAM (especially Erlang) programmers, there's a lot of conflicting information.
From my limited understanding, Kubernetes is comparatively complicated, and can hamstring BEAM instances with port restrictions.
On the other hand, there's a rarely documented soft limit on communication between BEAM nodes (informally, circa 70 units, IIRC). Above this limit, you have to make plans based on sub-clusters of nodes, though I have certainly not worked at that level of complexity.
Would be interesting to hear what other people think about this specific subject.
- di4na 6y agoThis limitation is far higher for the past few years and there is more work from the OTP team to raise it. Also there all the hooks needed to adapt and change it as you grow. I would argue that when you reach above a hundred nodes, you will need to optimise yourself in any tech though
- josevalim 6y agoFWIW, the soft limit for 70 units is about using the global module, which provides consistent state in the cluster (i.e. when you do a change, it tries to change all nodes at once). The default Erlang Distribution was shown to scale up to ~300 nodes. After that, using sub-clusters is the way to go and relatively easy to setup: it is a matter of setting the "-connect_all false" flag and calling Node.connect/2 based on the data being fed by a service discovery tool (etcd, k8s, aws config, etc). PS: this data came from a paper. I am struggling to find it right now but I will edit this once/if I do.
- strmpnk 6y agoI can also mirror this general guideline. I've run 250+ node erlang clusters just fine in the past. There were some caveats to how some built-in OTP libraries behaved but they were easy to replace or workaround in the past. That was many years ago as well. The distributed erlang story has improved with more recent releases (better performance on remote monitors for example) which might push the number a little higher than 300 if you are careful. Keep in mind the default style of clustering is fully connected so there is some danger in managing that many network connections (quadratically scaling for each node added) during network partitions which can be a problem if you're not tuning things like TCP's TIME_WAIT for local networking conditions. Even better, these days there are great libraries like partisan (https://github.com/lasp-lang/partisan https://github.com/lasp-lang/partisan) which can scale to much larger cluster sizes and can be dropped in for most distributed erlang uses cases w/o much effort.
- toast0 6y agoI 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/5f1ef352f971b2efad3ceb4030e2367e8996f893/lib/kernel/src/pg2.erl#L286 https://github.com/erlang/otp/blob/5f1ef352f971b2efad3ceb403...
- ksec 6y ago>I think our big cluster had around 2000 nodes when I was working on it. Is there fairly recent? I thought WhatsApp was on FreeBSD with Powerful Node instead of Lots of Little Node? >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. Something the I was thinking of when reading POWER10 [1], what system and languages to use with a maximum of 15 Core x 16 Socket x SMT 8 in a single machine. That is 1920 Threads! [1] https://www.anandtech.com/show/15985/hot-chips-2020-live-blog-ibms-power10-processor-on-samsung-7nm-1000am-pt https://www.anandtech.com/show/15985/hot-chips-2020-live-blo...