Using up almost half of your pps every 30 seconds for cluster maintenance certainly seems like it's more than "not a big deal", no?
Using up almost half of your pps every 30 seconds for cluster maintenance certainly seems like it's more than "not a big deal", no?
If your switching fabric can only deal with 1Gbps, yes, you've used it halfway up with heartbeats. But if your network is 1x 48 port 1G switch and 44x 24 port 1G switches, you won't bottleneck on heartbeats, because that spine switch should be able to simultaneously send and receive at line rate on all ports whicj is plenty of bandwidth. You might well bottleneck on other transmissions, but the nice thing about dist heartbeats is on a connection, each node is sending heartbeats on a timer and will close the connection if it doesn't see a heartbeat in some timeframe; it's a requirement for progress, it's not a requirement for a timely response, so you can end up with epic round trip times for net_adm:ping ... I've seen on the order of an hour once over a long distance dist connection with an unexpected bandwidth constraint.
It would probably be a lot more comfortable if your spine switch was 10g and your node switches had a 10g uplink, and you may want to consider LACP and double up all the connections. You might also want to consider other topologies, but this is just an illustration.
I think you’re missing the fact that the heart beats will be combined with existing packets. Hence the quoted bit. If you’ve got 1000 nodes, they should be doing something with that network such that an extra 50 bytes (or so) every 30s would not be an issue.
I personally never operated anything above roughly 250 nodes, but that limit was mostly due to following the OP’s advice about paying attention to the configuration of each node in the cluster. In my case, fewer nodes with fancier and larger raid arrays ended up being a better scaling strategy.
The network saturation is just a necessary cost of running such a massive cluster.
I really have no idea what kind of system would require 1000 nodes, that couldn't be replaced by 100, 10x larger, nodes instead. And at that point, you should probably be thinking of ways to scale the network itself as well.
> The network saturation is just a necessary cost of running such a massive cluster.
I think this actually answers it perfectly.
1. If you are running 1K distributed nodes, you have to understand that means you have some overhead for running such a large cluster. No one is hand waving this away, it's just being acknowledged that this level of complexity has a cost.
2. If heartbeats are almost 50% of your pps, you are trying to use 1Gbe to run a 1K-node cluster. No one would do this in production and no one is claiming you should.
3. If your system can tolerate it, change the heartbeat interval to whatever you want.
4. Don't use distributed Erlang if you don't have to. Erlang/Elixir/Gleam work perfectly fine for non-distributed workloads as do most languages that can't distribute in the first place. But if you do need a distributed system, you are unlikely to find a better way to do it than the BEAM.
Basically, it seems you are taking issue with something that 1) is that way because that's how things work, and 2) is not how anyone would actually use it.
Fairly similar, but smaller numbers in 2012 http://www.erlang-factory.com/conference/SFBay2012/speakers/...
The 2013 presentation is focused on MMS which I don't remember if it was as impressive: http://www.erlang-factory.com/conference/SFBay2013/speakers/... (note that server side transcoding is from before end to end encryption)
I don't think there were similar presentations on Erlang in the large at WhatsApp after that. Big changes between 2014 and 2019 (when I left) were
a) chat servers started doing a lot more, and clients per server went down on the big SoftLayer boxes
b) hosting moved from SoftLayer to Facebook and much smaller nodes --- also chat servers at SoftLayer were individually addressable, using (augmented) round robin DNS to select, at Facebook the chat servers did not have public addresses, instead everything comes in through load balancers
c) MMS was pretty much offloaded into a Facebook storage service (c++); not because the Erlang wasn't sufficient, but because MMS was loosely coupled with the rest of the service, Facebook had a nice enough storage service, a lot of storage, an awful lot of bandwidth, and it wasn't a lot of work for that team to also support WhatsApp's needs; also our Erlang MMS (and the PHP version before it) was built around storing files on specific, addressable nodes, but nodes at Facebook are much more ephemeral and not easy to directly address by clients
d) some amount of data storage moved off of mnesia into other Facebook data storage technology; again, not because mnesia wasn't sufficient, but more ephemeral nodes makes it cumbersome (addressable) and the available hardware nodes at FB didn't really match --- there's a very firm bias at FB towards using standard node sizes and the available standard nodes were like a web machine with not much ram or a big database machine with more ram and tons of fast storage; WA mnesia wants lots of ram but doesn't need a lot of storage (all data is in ram, and dumped + logged to disk) so there was a mismatch there --- things that stayed in mnesia needed much larger clusters to manage data size
Presentations became less common because of more layers to get approval, and also because it's less fun to share how we built something on top of proprietary layers that others don't really have access to. Anybody could have gotten dual 2690 servers at SoftLayer and run a nice Erlang cluster. Only a few people could run an even bigger chat cluster in a Facebook like hosting environment.
https://www.youtube.com/watch?v=A5bLRH-PoMY
40,000 erlang nodes in a cluster
lots of specifics about what tweaks they use. Rewatching, it seems like you don't have to really use the modified BEAM except in a few small soft-code (replace OTP/stdlib functionality) where they provide enough information that you could probably write it yourself if you had to, and a lot of the optimizations have been upstreamed into core BEAM -- WA is a pretty good citizen of the ecosystem.
Nodes only get so big. Way back when, a quad xeon 4650 v2 was definitely not 2x the throughput of a dual xeon 2690 v2; so you end up with two dual socket systems instead of one quad socket. A quad socket server often costs significantly more than two dual socket servers, and is likely to take longer between order and delivery. There's usually a point where you can still scale up, but scaling out is a better use of resources.