How Discord Scaled Elixir to 5M Concurrent Users (2017)
blog.discordapp.com
blog.discordapp.com
The Fortnite Official server has exceeded 100,000 concurrent users, Discord itself is way past that 5M concurrent number, we're now using Rust in certain places to make Elixir go faster, we've built a general purpose replacement to Process.monitor that scales a whole truckload more that we're open sourcing next week at Code BEAM SF... the list goes on.
There's a lot of fun stuff going on to try to make this system even more efficient and reliable, there's a lot to do still. We run everything on a very small engineering team (there are 4 fulltime engineers on the core infrastructure, only about 40 engineers in the whole company) and we're always looking for a few more. Feel free to reach out to me (zorkian#0001 on Discord) if this blog post sounds up your alley!
http://erlang.org/doc/tutorial/nif.html
> As a NIF library is dynamically linked into the emulator process, this is the fastest way of calling C-code from Erlang (alongside port drivers). Calling NIFs requires no context switches. But it is also the least safe, because a crash in a NIF brings the emulator down too.
Sounds like a pretty good use for Rust!
My only complain is the design of your blog. The header and the footer (to subscribe on Medium) takes so much space that there is relatively less space to read to actual content.
https://addons.mozilla.org/en-US/firefox/addon/make-medium-r...
That seems both risky and somewhat overkill considering its features. Does Firefox not support targeting a specific domain yet? Or is part of the problem that medium allows custom domains (it does, right?).
I can see in the extension source (thanks to https://addons.mozilla.org/en-US/firefox/addon/crxviewer/) that on every page, the extension uses JavaScript to check for a top nav bar or a login nag popup and hide them if present, then applies CSS that hides five other UI elements if they are present.
> Choosing to use and getting familiar with Erlang and Elixir has proven to be a great experience.
What background did the core infrastructure engineers have before tackling Discord in Elixir and Erlang?
Would this have been avoided have you started with golang or the JVM ?
We've since been able to squeeze out the performance from ETS we needed to be able to run our hashring. More on that in this PR here: https://github.com/discordapp/ex_hash_ring/pull/1
Also one amazing feat we've accomplished thanks to ZenMonitor (which we will open source soon as my colleague mentioned), is that we can now tolerate a guilds node failure and recover from it in ~40 seconds. We had a node fail on Friday actually due to the underlying host rebooting (gcp calls this a 'host error'), and didn't even notice until after the system had already recovered and the on-call got alerted after the fact. Back in the day this lead to a cascading failure throughout the system. A guild node runs between 600k-700k concurrent discord guilds, or servers as it's known in userland. And although we haven't done it recently, we clock a full restart of our distributed system at roughly 17 minutes (from shutdown to service fully restored).
Yeah, the most straightforward failure that we can’t migrate away from is the NIC (or rack switch) failing. Obviously, if your path to get off the box is dead, that’s not going to happen :).
Others though, like the hypervisor crashing or even host kernel are also possible, but much less frequent than “Hmm, I think the network is dead, we should start a replacement VM”.
One of the many motivations for (now) only having Persistent Disk for boot disks, is that it lets us avoid a whole class of truly unrecoverable errors. There are still bad DIMMs, but monitoring for ECC failures often lets us mark the host for replacement in time for the VMs to migrate off before they actually fail.
As a mostly Elixir programmer for 2 years now, I find it quite amusing how the Kubernetes community tries their damnest to emulate Erlang's OTP system -- going as far as to bolt stricter typing system on Golang even -- and try and reinvent Erlang's "let it crash and get rebooted" idea.
I am however not at all convinced that "let it crash" is a good mantra when applied to entire containers. One such container can take 20+ seconds to restart and you can lose a lot more compared to the mini-processes Erlang/Elixir have which are happily left to crash and [semi-]auto-recover. Imagine if that container was ingesting events and crashed when it had 5_000 in its in-memory queue and only managed to process 50-100 of them. Or imagine 100 transactions in progress being cancelled.
I applaud the hard work of the Kubernetes team, I am just not very sure they invest their energy in the right tech stack. Sure Golang is faster than Erlang/Elixir -- by a lot, too. But is it not K8s idea to be fault-tolerant and not immensely quick? What does a Golang's speed matter when a container needs seconds to reboot? In my eyes, K8s could have been written in Bash shell scripts... but I am probably missing something important.
Go makes sense as an implementation language for a system like this for many reasons:
- it doesn't require a VM (goodbye Java) - it is a memory safe language (goodbye c++) - it has strong types to improve reliability - it does a good job of handling complex, multi-module code efficiently (these two say goodbye to python)
So of the Google friendly languages you end up with Go.
And I don't know if it was planned, but Go has turned out to be a boon for contribution. A surprising number of devops engineers are willing to dip their toes in Go waters and I highly doubt they would've done so with Erlang.
I think you'd end up going the chef route, writing the control plane in Erlang and expecting users to interact in Ruby or something like it.
Also this whole conversation seems a bit off because Kubernetes is a multi-service architecture with many components. You can absolutely write operators in other languages since Kubernetes exposes an API.
It's kind of the whole point to be able to take an existing application and run it in k8s rather than a regular VM with minimal changes. Expecting developers to rewrite everything in Erlang is nuts.
...But I never said that?
You are correct on your points and I don't disagree. Golang is certainly a much better choice than Bash indeed. You are also correct on developer willingness to work with Golang.
I am aware that K8s and Erlang/OTP are apples to oranges comparison; they serve different needs. Whereas Erlang's runtime (the BEAM VM) can give you fault-tolerance, K8s tries to do roughly the same on a higher level - it tries to give you throwaway containers that can be switched off and on at any time.
As I mentioned in my parent comment, I admire their work.
What my point was that if you squint your eyes hard, it kind of looks like K8s wants to invent Erlang/OTP for infrastructure (as opposed to Erlang/OTP which gives its guarantees per node).
There is no such thing, because Nodes (and even entire datacenters) can fail.
In other words, node failure is an infrastructure problem that is best NOT handled by your bespoke application code. Replacing failed nodes should NOT be custom code in your app, that way lies madness.
> kind of looks like K8s wants to invent Erlang/OTP for infrastructure
You can say "K8s and Erlang implement similar ideas implemented at different levels." But you can't pretend one is a substitute for the other, nor that "Erlang has done everything that K8s can do."
E.g. Writing in Erlang doesn't magically get you:
- deploy/upgrade of the Erlang runtime, including rollback + multiple versions co-existing, including any compiled foreign code (which is required in the article). - ability to log, monitor, probe and reroute the connections between services -- In a standard way such that the application doesn't have to be modified.("service mesh") - Ability for an entire ecosystem of tools to inspect the versions of your application services that are deployed, because it's exposed as an API. - A standard ecosystem of plug-ins for operators (autoscale, autoscale to EC2 spot instances, capacity planning, "best practices" of running a MySQL cluster, etc.) None of these should ever be mixed with the application. (Unless you are Kelsey Hightower: https://github.com/kelseyhightower/hello-universe )
Gross over-simplification that I'd also call a strawman. In the face of lack of electricity, of course no computer language matters at all. What's your point?
Neither does Java, only devs clueless about Java world aren't aware of AOT compilers to native code.
Also sorry to spoil your fun, Kubernetes was originally written in Java, it was later rewritten in Go as other team took over.
Your comment about my experience with Java was needlessly provocative so I'm going to ignore it.
Nowhere in my commented I addressed the Java remark to your specific person.
It was targeted to all of those that bash Java without actually knowing how the rich the eco-system actually is.
On my previous comment I explicitly mention it was targeted to developers that bash Java, without having any clue what they are bashing about, cargo cult if you wish.
So either you feel offended, because you are indeed attacking Java without having any clue about the Java eco-system, or are deciding to play victim here, when I never mentioned you directly.
This is a public forum, so I leave to everyone else and HN moderators to judge my comments, and won't reply to anything else on this thread.
> I think you'd end up going the chef route, writing the control plane in Erlang and expecting users to interact in Ruby or something like it.
DevOps engineers have learned to be comfortable with Ruby (Chef) and Python (Ansible) so this certainly makes sense.
OTP doesn't do 5% of what Kubernetes is able to provide ( cpu / io / memory quota, live / readiness probes, rolling deployments / canary / blue green.
I am well aware of the differences in scope between K8s and Erlang/OTP. My point was that they kind of try to do the same only on different levels.
And I am still not sold on the idea of throwaway containers. That works well if you have a hyper-network of microservices that are able to discover each other and self-heal a bigger graph of services but if you have any sort of a more classic coherent whole app... then not really. Throwing containers away and rebooting new copies isn't something that looks viable for many projects.
But hell, who knows. K8s team and their audience are really dedicated. They might change every single tool of the trade just to make K8s work well. I was mostly saying that I don't feel it brings something really radical to the table.
Caveat: in the environment I work, we don't do real gen_server:call, because all of the included monitor/demonitor calls are too expensive for us. The trade-off is we only have timeouts when the server goes away during the request. If you had the monitor letting you know the request crashed the server, you could presumably do a smartish retry -- but it's still nicer to only need to do that for the crashing request, not the others that would fail simply because the server died.
If you want to re-invent the whole Erlang eco-system you should just switch and call it a day, the chances that a company that needs to get a job done will be able to pull this off successfully as a side project are nil.
Erlang/Elixir are still quite impressive out of the box, even without these Discord-specific optimizations. But they probably wouldn't scale to 5M right away.
The real win IMO is the good reliability:performance ratio that Discord achieved. I feel Erlang/Elixir are excellent in optimizing this exact metric.
But it's still cool to see it done, and a lot of people are unfamiliar with Erlang or Elixir, so the discussion is interesting.
I think it’s probably healthier to cheer the people doing it right than to keep shaming everyone else.
This is really surprising to me, and definitely something the elixir team should look at optimizing. Sending messages should be extremely fast.
In my implementation, only 1/2 a million concurrent IOT devices, I use routing tables that narrow down a process to a node and registry cluster (https://hexdocs.pm/elixir/master/Registry.html) , and from that registry it fans out to 1 of 100 supervisors for that worker type per node.
SocketCluster has been used in production to service hundreds of thousands of concurrent users and it can handle millions.
I've had one report of a chat system (adult industry) which could handle 250K concurrent users using only two large servers.
Also, I once did some consulting work for a popular cryptocurrency trading platform which handled tens of thousands of concurrent users/trading bots (with very high frequency of messages). After I was done, that company didn't talk to me for 6 months straight; it turns out that they hadn't had any issues with their pub/sub cluster since.
Unfortunately SocketCluster doesn't get discussed very often among influential circles. I have no idea why because the feedback I get from users is essentially 100% positive.
I guess Node.js doesn't get much hype these days.
I still love NodeJS ! but once you learn Erlang its hard to go back, it reminds me of PG's having a higher bird eye view.
NodeJS, Python, Haskell (!), PHP, C - all belong to the same class of coding style with the same type of problems.
You use Erlang not because of playing the testosterone game of nominal performance - but because you really want some guarantee.
( Money should be not a problem since the world seems flush with cash - if your manager is complaining its because he wants his bonus to be higher. )
Think about how prinf / console.log / print .... works in traditional settings.
- How would you make it so that printf doesn't crash your entire program if the console hangs.
- how would you isolate a single codebase's IO operations ?
- How would writing to console work in a multi threaded environment ? multi server ? 100 servers ?
Haskell sort of tries to answer these questions but I am not sure how its going go work out for them, in Erlang's process based universe all those questions have been answered already !
I do really appreciate Erlang / Elixir (have contributed several libraries), but the problems you describe are not uniquely solved by Erlang. Akka / Scala is another take on the whole actor based architecture, and seems to have considerably more traction (hiring talent will be easier for your manager).
They solve a subset of the problems that Erlang's OTP solves. They don't have the entire package.
I agree, I do not any Scala / Akka experience so I cannot argue for or against due to ignorance.
But as you say at least the platform / language addresses these concerns.
Why would the console hang? If that happened, it would signal a major issue at the OS level and probably not related to your application (unless you're trying to log an extremely massive string; which is a bad idea and you'd probably already have run out of memory before that could happen). I have never seen the console/stdout hanging in production and I've built some pretty high-traffic distributed systems with Node.js.
>> how would you isolate a single codebase's IO operations ?
What sort of IO operations are we talking about? Network, Disk? There are many ways to inspect different kinds of IO operations. The Node.js ecosystem offers a large number of modules which would let you achieve that.
>> How would writing to console work in a multi threaded environment ? multi server ? 100 servers ?
Node.js is perfect for running on Kubernetes. There are many K8s dashboards and tools which allow you to browse and aggregate logs from thousands of machines with very little effort. I don't see how this point has anything to do with Erlang specifically. A language-agnostic container orchestrator like Kubernetes is the best way to go over a tool which only works with a specific language.
There are Node.js frameworks which offer kubetnetes .yaml files and CLI tools which allow you to deploy a highly scalable cluster to K8s in a few minutes.
By default (and unless you go out of your way w/ web-workers) the javascript event-loop (and thus nodes event loop) is single threaded. To work around this, you can bind many node processes to a given port to load balance requests (SO_REUSEADDR, anyone?) - or simply run many smaller instances of node (perhaps in a bunch of containers) where traffic ingresses in via some form of load balancer. The load balancing problem is unavoidable, and you will definitely need the same if you want to send requests to multiple BEAM nodes. However, BEAM can schedule your work across all the cores you give it.
But let's talk about work for a second, and about a very special thing that the BEAM VM gives you, that other runtimes (whether it be node, JVM, golang's, etc...) aside from the actual operating system of your computer does not. And that's specifically preemptive scheduling.
Suppose you have a single core computer that's running Linux, and you have a process that is sitting there busy looping. Let's say we just make a simple script that does nothing infinitely in a loop. Does your computer grind to a halt? Most likely, no. You can still probably use your terminal, move your mouse, operate your web browser, etc... You can thank pre-emptive scheduling for that. The OS suspends the process to allow other processes to do work - hopefully in a fair manner (on linux, the CFS (aptly named Completely Fair Scheduler) does this).
Now let's say you have a single node process serving requests. Let's say that a specific kind of requests requires 250ms of CPU time to compute - and does not explicitly yield back to the event loop. (You can imagine doing some processing of input data, deserialization, serialization, aggregation, etc...). During this computation, nothing else within the node process can progress. This means that requests that may not take a lot of time to compute now have to wait 250ms to be processed. Generally, I see node deployments not having single request/response request handling, but rather many concurrent requests/responses being handled at any given time, using promises/callbacks to allow the event loop to progress while waiting on IO from something else. During the periods of expensive computation from a given request handler, the entire event loop is stalled, and the response time percentiles of your requests spike. A pathological case would be something like `setTimeout(() => while(1) { }, 1000)` deadlocking your entire node process after a whole second, as the loop does not yield back to the event loop.
In BEAM, this does not exist. Processes are scheduled and pre-empted - to allow for fair utilization of the underlying computation resources (very much like how your OS does it.) This means that a computationally intensive process does not stall the event loop for all other processes, meaning that your response times and percentiles remain low for all other work within the system.
Now of course, you could hand-craft your javascript code to explicitly yield to the scheduler every so often, but that's a lot of work that you as a programmer are now doing that your runtime could be doing for you, and if you forget to do it, could be catastrophic to the performance of your soft-realtime system.
This is only one of the many benefits that OTP/BEAM provide over other runtimes. But one compelling enough for Discord as a company to bet on it. For a given service, we run entirely homogeneous infrastructure. We do not need to allocate or dedicate special resources to our largest servers (100k CCU/350k members), and instead can run and schedule it alongside the millions of other small servers that exist on Discord - all without negatively impacting the performance, percentiles, or soft-realtime guarantees of your chat with a few of your friends.
In any case, you don't necessarily need loadbalancing at the host level, you can load balance at the cluster level only. Your load balancers (e. g. nginx or haproxy...) have their own hosts/machines in your cluster and they loadbalance between processes directly. Some of those may be running on the same machine but the load balancer does not dostinguish between them. A random load balancing approach yields the most even distribution from my experience - You do need each process to be able to support maybe 1k concurrent users in order to get the sample sizes on each prpcess to allow even random distribution between them but this is easily achieved with most Node.js WebSocket libraries. They can easily support 10K concurrencr connections per prpcess with very high message throughput. If the commections are mostly idle, each process can handle 100k connections or more.
We use BEAM/OTP for way more than just holding open websocket connections. Our entire websocket layer is a few hundred lines of elixir code - and honestly hasn't been touched in over a year - and has remained pretty much the same as we scaled from 200k ccu -> to well over 5m ccu. Holding open websockets and load-balancing them is pretty much a solved problem for us.
TL;DR: No, no other language or framework in the world has the primitives that Erlang and Elixir have.
People on HN and Reddit really love acting non-impressed and claiming the pain points are easily solved in other languages.
My 17 years of career say this is not true at all.
Issues I have raised are 100x harder to address then some syntax.
I mean supposing we can get past console printing issues?
- For example I have a project involving smart electrical inverters, I need some guarantees regarding crash handling / low latency. Not a lot of scaling issues.
- With scaling involving NodeJS, I consistency had issues with exploding RAM and crashes - unable to isolate part of codebase.
I struggled a lot to fix these problems, so while searching for a solution I came across Erlang and haven't looked back.
Design for resilience is not a small thing. I suggest to read Joe Armstrong's "Making reliable distributed systems in the presence of software errors" [1] as a starting point on this.
"Pet vs Cattle"
The thing I don't like about Erlang is about the runtime that is mixed between code and infra which I think is not a good idea, it was designed before we made progress with HA platform like Kubernetes. https://github.com/kubernetes/community/blob/master/sig-scal...
It should be separated and it's what pretty much everyone is doing nowdays.
If you rely on the infra for recovery, you’re going to be in serious trouble. There may be 300k users connected to an app instance, and you’re going to be kicking all of them out and restarting an instance every time the minor feature is called. This turns a minor bug into a full-blown outage.
Now we’re talking about a “printf crash” as an example, but in my experience the issues are more subtle. I’ve seen a Python app leak file descriptors due to a bug in an object’s destructor, which only happened when an exception was thrown in a certain place. This caused a service outage as the app looked “fine” from the infra’s point of view (it could respond to health check requests), but the functionality was dead. With Erlang processes, using a process per connection, the VM guarantees resources are cleaned up when a process dies, so that kind of issue doesn’t happen.
The Zen of Erlang goes a bit deeper into transient issues and why the supervision model helps: https://ferd.ca/the-zen-of-erlang.html
I actually think that Elixir / Erlang are actually not great for those kind of problems because they're slow language and consume a lot of memory. They allow you to do easy message passing + horizontal scaling but the runtime is inefficient. Java / C# net core / Rust / C++ / Go are much faster than Elixir. ( if you actually read on it you'll see that they use a lot of C / C++ / Rust to make it fast which is not something you need to with the above languages ).
And for deployment / scaling just use Kubernetes or equivalent, better than BEAM trust me.
The power of BEAM is that although the performance may not be the best, throughput and response time is consistently low throughout the system - allowing no single process (or actor) from monopolizing resources of the system. When you use a homogenous server configuration like we do, this makes a lot of sense. Our largest guilds (100k ccu, 350k members) are on the same nodes as all of our other servers. And when they're busy, the small guilds notice no performance degradation.
Other chat products out there (that I hear use the JVM for their real time stuff) have to spin up dedicated resources to hosting their larger servers/clients - and even then cannot handle servers with as many users or concurrents that we can. Every single discord server runs in a homogenous cluster, without special dedicated resources for our largest instances.
Could you just get rid of red messages and make it "pending" until acked as "sent"?
https://www.bignerdranch.com/blog/elixir-and-io-lists-part-2...
Message passing is somehow slow as well, slower than go(Lang) of you want a comparison.
But on those systems, the one that you design on the BEAM, velocity is usually not a problem. What you usually try to do is to reach a design that does not have a single bottle neck or failure point so that pretty much whatever happen you can just add machines.
This turn out to be a great way to design multi{processes, cores} and multi nodes system.
What you usually get on a BEAM under load is an 100% of CPU usage, even on multi core, but to be honest those core are not used as efficiently as you could.
It is a matter of tradeoffs, I can quickly write multi process software that can easily scale, but I will leave on the table some raw performance number.
Again it turns out that just raw numbers are not as important as they are simple to measure.
https://www.bignerdranch.com/blog/elixir-and-io-lists-part-2...
- Reading from the mailbox?
- Destructuring and binding the content?
- Operating on the data by transforming, modifying, or filtering it?
etc.
More interesting might be to expose mnesia, maybe. But even then, it might be most useful to build out your data handling in Erlang and expose a higher level API for clients in other languages.
Discord already uses Cloudflare, I'm curious what they think of Workers.
edit: I'm on macOS, not an heavy user, I used Slack / Skype / Hangouts / WhatsApp / Messenger / etc. before
- Abusive trust and safety team - Constant outages for both users and bots - Frontend is heavily bloated - UI designed for money grab instead for the users
While I appreciate their technical accomplishments, this is one of those things that gives me chills about the company.
Don't fall in love with services. (Or with people ;P)
Not to mention where an update randomly the noise suppression would completely mute your mic.
We are a very heavy user of cloudflare workers. Our marketing/developer/web app are entirely served from the edge using workers - and we use workers for request authentication for downloading game chunks on our store.
Discord has actually been using workers since before they were available to the general public. They've been using them mainly for edge caching games for the store, A/B, and promotion of clients to different testing environments.
Never in the world would I have thought a gaming chat client had more robust features than the pricey enterprise chat client
It's the replacement for the vintage volume control, while a little sub-par by any modern standard, you can pick the input/output device per application
No offence but that's like your opinion.
Coming from IRC and TeamSpeak, to me Discord is unnecessarily bloated and full of analytics.
>I used Slack / Skype / Hangouts / WhatsApp / Messenger / etc. before
Yeah, explains it. :P
The voice service isn't very high quality and there are frequent server outages.
I use Telegram, Facebook, whatever except Discord for instant messaging, and wish Discord wasn't the default for non-professional groups.
People my age (upper 20s) praising it are always those who used Skype for gaming and never touched Teamspeak/Ventrilo/Mumble.
As gaming itself got broader reach a lot of casual folks were left to their own devices and underserved, so that’s the big gap being seen. I know plenty of 20-somethings that know about and used the predecessors to Discord but they’re all hardcore gamers compared to even the folks closer to mid/late 30s that may have been hardcore before but don’t have the time to fiddle with these systems anymore.
It seems with Discord quality was tough to maintain at scale and we’ve got the inverse problem set.
Now, you might want a bouncer to store messaged and logs when offline, but not to be able to connect to multiple servers. Logging is a basic feature in IRC.
And it is, of course, closed source, so no way to use an alternative client. I wish I could use an IRC gateway for it, that would be cool (I very, very rarely use the voice chat)
Just because something is closed source, doesn't mean it's impossible to reverse engineer. The entire system is compiled into a javascript webpack, alongside just reading web requests/their documentation, doing some basic functionality in an alternative implementation is really easy.
This apparently breaks the terms of service, but they are only enforce this if you're an outlier for the number of API requests you send.
The most well known 3rd party client would be Ripcord, which I use for it's Slack features
This thing usually goes: Until the better client threatens the proprietary one with is better monetized.
Then I know they're reading.
>and you have no idea who might be saving logs or where they might be publishing them
If I'm in a room with you nothing stops you from telling others what I've told you.
Sure, but I won't have a complete written record like an IRC log, and won't be able to credibly quote.
With how discord uses semaphores, the Consumer will not even make a demand for more events from the producer.
Other than that, I think in the case of Discord, it's more of an RPC thing, and GenStage is rather made for concurrent streaming and processing of data.