13 ms·
How Discord handles over a million requests per minute with Elixir’s GenStage
- bpicolo 10y agoI love Discord, and love Elixir too, so this is a pretty great post. Unfortunate that the final bottleneck was an upstream provider, though it's good that they documented rate limits. I feel like my last attempt to find documented rate limits for GCM/APNS was fruitless, perhaps Firebase messaging has improved that?
- hotdogs 10y ago"Obviously a few notifications were dropped. If a few notifications weren’t dropped, the system may never have recovered, or the Push Collector might have fallen over." How many is a few? It looks like the buffer reaches about 50k, does a few mean literally in the single digits or 100s?
- Sikul 10y agoGood question. We don't have metrics on the exact number dropped. We're using an earlier version of GenStage that doesn't give any information about dropped events. Once we upgrade we'll have a better idea.
- teej 10y agoThat seems too important to have zero visibility on to me. Just eyeing the graphs, your queue size grew at 750m/s from 17:49 to 17:50. You then starting shedding at 17:50 for 40s. Assuming the ingress rate was roughly linear (which it looks like it was) you shed ~30,000 requests out of 3-4M. Does that not seem high to you? This system seems great for at most once delivery. I wish I had more problems to solve with that constraint.
- Sikul 10y agoYep, such is life with new tech. GenStage was 0.3.0 when we started using it. We'll be able to get this visibility once we update to a newer version (hasn't been prioritized). FWIW, the buffer only fills up to the peak about once a month. Load shedding is the last ditch effort in catastrophic situations so that everything doesn't fall apart.
- bcherny 10y agoCan you explain why it's necessary that some notifications were dropped?
- Vishnevskiy 10y ago2 reasons. - To avoid OOMing the Erlang VM. - If the notification queue is backed up then older notifications are not worth delivering if we can speed up delivering of more recent ones.
- DougN7 10y agoI was wondering the same thing. Dropping an unknown number of requests isn't all that impressive. It seems like a simpler approach would have been to use a Message Queue of some sort with pushers pulling items from the queue.
- poorman 10y agoThat's awesome and it just goes to show how simple something can be that would otherwise involve a certain degree of concurrent (and distributed) programming. GenStage has a lot of uses at scale. Even more so is going to be GenStage Flow (https://hexdocs.pm/gen_stage/Experimental.Flow.html https://hexdocs.pm/gen_stage/Experimental.Flow.html). It will be a game changer for a lot of developers.
- pwf 10y ago50k seems like a low bar to start losing messages at. If this was done with Celery and a decently sized RabbitMQ box, I would expect it to get into the millions before problems started happening.
- ramchip 10y agoAt 15k notifications per minute, a million notifications would take 1hr to clear before the queue returns to normal. I would imagine they prefer to shed load early so notifications don't get delayed, hence the small buffer.
- jhgg 10y agoAt some point, when a system has entered a failure mode for a while, it makes sense to start shedding load, rather than attempting to deliver every single push notification. Also worth mentioning, a minute of downtime is already a million backed up pushes. Beyond that, it becomes infeasible to attempt deliver them. Edit: Also worth mentioning, the 50k buffer is for a single server, we run multiple push servers in the cluster.
- Vishnevskiy 10y agoThese machines do more than just push. They also buffer messages for each individual user to "potentially" push if they don't read them on the desktop client. This happens before the flow this article talks about. We currently have 3 machines doing this for millions of concurrent users. At the writing of this article it was 2 machines.
- jsjohnst 10y agoWhat size machines are these? I'm shocked that this volume is your max handling with Erlang unless your using a smaller T series AWS instance for this.
- Vishnevskiy 10y agoThese are n1-standard8 on GCE. These are getting easily over 30,000 requests a second each about updating queues for new messages. And also are subscribed to presence events from our presence system to millions of people. It is a very busy service ensuring we only deliver messages to people not at their computer.
- AgentK20 10y agoAnyone know of a equivalent libraries like GenStage for other languages? (Java, NodeJS, etc) I'd definitely be able to put to use things like flow limiters and queuing and such, but none of my company's projects use Elixir :(
- bhelx 10y agoAkka streams?
- bpicolo 10y agoReactiveX seems to have documented notions for it: https://github.com/ReactiveX/RxJava/wiki/Backpressure https://github.com/ReactiveX/RxJava/wiki/Backpressure Highly recommend the Reactive series of libs. They're typically very well done. The guy below is right that Akka is perfectly suited.
- gazarullz 10y agoFor java there's also: - Project reactor from Spring - Reactive Spring (following up with spring 5.0)
- wtf_is_up 10y agoThere was an initiative not long ago called Reactive Streams which established some common interfaces to build things like this. Back pressure was one of the main concerns. Some implementations are listed here: http://www.reactive-streams.org/announce-1.0.0 http://www.reactive-streams.org/announce-1.0.0
- jtchang 10y agoThe most important part of this article is the concept of back pressure and being able to detect it. It's common in a ton of other engineering disciplines but especially important when designing fault tolerant or load balancing systems at scale. Basically it is just some type of feedback so that you don't overload subsystems. One of the most common failure modes I see in load balanced systems is when one box goes down the others try to compensate for the additional load. But there is nothing that tells the system overall "hey there is less capacity now because we lost a box". So you overwhelm all the other boxes and then you get this crazy cascade of failures.
- user5994461 10y agoYes. You need to adapt the capacity of the system to handle the full load with -N- boxes dead. Corollary: If you have 2 boxes, each of them has to be able to handle all the traffic, so you can't save money by using smaller boxes :D Corollary #2: If you have 2 datacenters, each of them has to be able to handle all the traffic, so you burn a lot of money :D
- xxpor 10y agoOr just the ability to tell your clients to back off. Hopefully your clients retry server errors with exponential backoff, if you lose a datacenter you can send half of the requests 503 until you're back at a manageable load. Hopefully detecting the load/generating the 503 is really cheap.
- user5994461 10y ago> until you're back at a manageable load. Nope. Gotta have the capacity at hand. You can't rely on getting it later when everything is fucked and you're already down. Down = Loosing clients = Loosing money = NoNonononono Source: Lost 4 million dollars last week because of that. Exponential backoff is to prevent cascading failure (and retry causing DDoS), that is a failsafe in case of failure, but an excuse for failing.
- 10y ago
- dimino 10y agoWhat is up with Discord? I feel like it's quietly (maybe not so quietly) one of the bigger startups to come out in the last two years. It seems to have totally taken over a space that wasn't even clearly defined before they got there.
- Numberwang 10y agoI ended up there for the first time last night and must say that there is a lot to like about it. I found some good communities and integrating media and so on all felt quite streamlined, and the system was snappy. It's just too bad there are a dozen IMs/voice/video and a dozen slacky/feed companies.
- bpicolo 10y agoThey had a really well defined user-space, marketed at it well, and really nailed the user experience, while still being free for the typical user. There is a lot to love about Discord.
- baldfat 10y agoI use Discord everyday BUT I seriously prefer IRC with weechat and glowing-bear.org. I feel like everything down with Discord could be done with IRC in a open source way. IRC for the 21st Century?
- bpicolo 10y agoIRC is a pretty not-extensible platform. It's just not in the spec to have larger messages, metadata. Then you get into things like push notifications, image hosting, video hosting, file hosting... IRC for the 21st century IS these apps like slack, discord, etc.
- superkuh 10y agoIRC is for people who participate in the internet. For them it isn't hard to host their own website or just use any of the innumberable image/video/whatever hosts and paste a link. Discord is for tech ignorant people who consume things pushed to them by companies over the internet. Bundling up everything into a centralized, proprietary package will go bad for the consumers eventually. But for now they'll trade away their freedom for convenience and enjoy it.
- erikbern 10y ago"requests per minute" is such a useless unit of measurement. Please always quote request rates per second (i.e. Hz). Makes me think of the Abraham Simpson quote: "My car gets 40 rods to the hogshead and that's the way I likes it!"
- ceejayoz 10y agoI, like most people, have no idea what a rod or a hogshead is. The same is hardly true for the conversion of minutes to seconds.
- StavrosK 10y agoSuch arcane units as "seconds" are only used by three countries in the world, though.
- hueving 10y agoHere's a cool trick I figured out. If you have something measured in units per minute, you can divide it by 60 to get units per second. I won't even charge you to use the method even though I'm in the process of patenting it.
- corobo 10y agoYou can also do it the opposite way if you want less specific numbers. Multiply by 60 and round off for the units per hour!
- deleted 10y ago[deleted]
- user5994461 10y agoActually. The conversion doesn't work. The requests per minute number is an average. The requests per second number should be given for peak load. That is a very important metric, a system has to be scaled to sustain the peaks, not the average. We'd need to know the traffic pattern to know the multiplier, that is certainly not 60 :p
- coverband 10y agoQuick serious question: How does this company plan to make money? They're surely well funded[1], but what's their end game? [1] "We've raised over $30,000,000 from top VCs in the valley like Greylock, Benchmark, and Tencent. In other words, we’ll be around for a while."
- meddlepal 10y agoThey have an awful lot of information about video gamers in conversation history. They could mine that data for game companies and sell it as a way to help companies build better, more addictive and mechanically pleasing games.
- 0942v8653 10y agoAggregated or Non-identifiable Data: We may also share aggregated or non-personally identifiable information with our partners or others for business purposes.
- falcolas 10y agoPure opinion: > mine that data [...] more addictive [...] games Oh, fuck no. No, no, no. I can not think of a more abusive thing for a company to do to its customers than that suggestion right there. How little respect would a company have for their fellow humans that anyone could even consider such a move? Maybe that's just a failure of imagination on my part, but I'm ok with that.
- chinhodado 10y agoI'm not sure how useful that is. It's not like gamers' opinion about games are hard to come by. Gamers are very vocal about their opinions, so any game developer looking for feedback can just go to Steam/Reddit/NeoGAF/whatever and read to their heart's content, or maybe even communicate directly to their player base.
- KMag 10y agoWhat people say/believe they like isn't necessarily consistent with their behavior. Data-mining with sentiment analysis might be useful. ("'Wows' dropped 1.3% and negative-sentiment profanity rose 0.5% after the latest patch.")
- sbov 10y agoIs the number of Push Collectors to Pushers constant or can it vary based upon notification load?
- jhgg 10y agoIt is constant - but iirc, it'd be trivial to make a dynamically scaling pool. At the end of the day, a pusher is just a TCP connection. Keeping a pool of fixed size and planning capacity around scaling horizontally is a perfectly acceptable approach - given you know the potential throughput for each pusher.
- mevile 10y agoI spend a lot of time in the PCMR Discord, which is pretty lively. The technology seems to be solid, while the UI has issues (notifications from half a day ago are really hard to find for example on mobile devices). Otherwise I'm on Discord every day and love using the service. I miss some slack features, but the VOIP is very good.
- b1naryth1ef 10y agoWhat features in particular? The most common one we hear is search, which is actually implemented and undergoing internal testing before a public preview soon.
- mevile 10y agoIt's just what I mentioned. I'll get a notification, and I just can't find where I was notified from. Like on Android, if I click on the notification I would expect it to take me to where the conversation happened where I was notified. It would take a really long time of scrolling to try and find the notification given the volume of discussion that happens. Can I just like click on something to see all my notifications from android, click on them and go to the conversation?
- snambi 10y agomillion requests per minute, is this a big deal?
- user5994461 10y ago16k per second. 83k per second during peak (assuming 80/20 default traffic rule). - 100 /s = typical limit of a standard web application (python/ruby), per core - 1.000 /s = typical limit of an application running on a full system - 10.000 /s = typical limit for fast systems (load balancers, DB, haproxy, redis, tomcat...). - Over 10.000/s You gotta scale horizontally because a single box [shouldn't] can't take it. The difficulty depends on the architecture and what the application has to do (dunno, didn't go through the article). You make something that can scale by just adding more boxes, then it's trivial, just add more boxes. Well, it's gonna costs money and that's about it. So no. Not a big deal at all... if you've done that before and you've got the experience :D
- manigandham 10y agoWhile 1k/sec seems to be an average throughput for most web apps due to all the logic, 10k/sec is nowhere near the limit for fast systems, many can do well into 6 figures per second with some now doing millions/sec.
- user5994461 10y agoRight. 10k is not a hard limit. It's the standard I expect, for real world applications, on classic server hardware, with limited tweaking. The 6 figures benchmarks that send/receive 1 byte data with all unsafe flags enabled are not representative of real usage.
- manigandham 10y agoYour description for fast systems refers to high-performance software like load balancers and databases. In this case, 10k is nowhere near the limit on modern machines, they all do 6 figures per second.
- user5994461 10y agoI'd like to say that the official performance unit is the "request per second". And its cousin, the requests per second in peak. The average per minute only gets to be used because many systems have so little load that the number per second is negligible.
- imaginenore 10y ago> "Firebase requires that each XMPP connection has no more than 100 pending requests at a time. If you have 100 requests in flight, you must wait for Firebase to acknowledge a request before sending another." So... get 100 firebase accounts and blast them in parallel.
- rv11 10y agojust wondering, what is the difference if I use two kind of [producer, consumer] message queues (say rabbitmq) instead of this? Does genstage being a erlang system makes a difference?
- di4na 10y agoRabbitMQ is written in erlang. So basically you use it natively instead of bringing and configurating a big dependency. It just come with your language for free without needing another process, etc etc.
- sandGorgon 10y agohow does one achieve this in Celery 4? I remember there was a celery "batch" contrib module that allowed this kind of a batching behavior. But i dont see that in 4
- manigandham 10y agoAkka(.NET) or any actor system is a perfect fit for this and brings the same functionality to other languages and frameworks.
- brightball 10y agoNot exactly. Without running on the BEAM you're left with cooperative scheduling (handing back control to the scheduler) of processes instead of pre-scheduling (the scheduler can stop you). That makes it possible for one processor heavy operation to take over and slow down everything else. BEAM ensures that if you have millions of request coming through and suddenly 1000 4 day long operations kick off on the machine, that the millions of normal, smaller operations continue responding and performing as expected. Fairly critical for the stability of real time systems. The other piece here is that these processes are cheaper on the BEAM than any other platform in terms of RAM cost. .5Kb / process on the Erlang VM. A goroutine in Golang is the next closest at 2kb. The two combined are one of the big reasons why benchmarks don't tell the whole story with Erlang/Elixir. It's harder to measure consistency in the face of bad actors.
- IOT_Apprentice 10y agoWhy not use Kafka for back pressure?
- jondot 10y agoHate to be a party pooper, but I'd like to give people here a more generic mental tool to solve this problem. Ignoring Elixir and Erlang - when you discover you have a backpressure problem, that is - any kind of throttling - connections or req/sec, you need to immediately tell yourself "I need a queue", and more importantly "I need a queue that has a prefetch capabilities". Don't try to build this. Use something that's already solid. I've solved this problems 3 years ago, having 5M msg/minute pushed _reliably_ without loss of messages, and each of these messages were checked against a couple rules for assertion per user (to not bombard users with messages, when is the best time to push to a a user, etc.), so this adds complexity. Later approved messages were bundled into groups of a 1000, and passed on to GCM HTTP (today, Firebase/FCM). I've used Java and Storm and RabbitMQ to build a scalable, dynamic, streaming cluster of workers. You can also do this with Kafka but it'll be less transactional. After tackling this problem a couple times, I'm completely convinced Discord's solution is suboptimal. Sorry guys, I love what you do, and this article is a good nudge for Elixir. On the second time I've solved this, I've used XMPP. I knew there were risks, because essentially I'm moving from a stateless protocol to a stateful protocol. Eventually, it wasn't worth the effort and I kept using the old system.
- Vishnevskiy 10y agoI think you misunderstand the problem we are solving here. We are not trying to solve this because our system can't handle it. We are protecting it from when Firebase decides to slowdown in a way that causes data to backup and OOM the system. Since these are push notifications that have a time bound on usefulness we don't care about dumping to an external persisted queue like RabbitMQ or Kafka (we rather deliver newer notifications faster, than wait for the backed up buffer to flush). Firebase also only allows 1000 concurrent connections per senderId with 100 inflight pushes (that have not received an ack) which means that only 100,000 can be inflight. Ultimately if a remote service is providing backpressure because it is having a struggle no amount of auto scaling on your end is going to help you. This service buffers potential pushes for all users being messages, that then watches the presence system to determine if they are on their desktop or mobile (this is millions of presence watchers and 10s of millions of buffered messages), and users are constantly clearing these buffers by reading on the clients and finally when a user is offline or goes offline we emit their pushes to them (which is what this article talks about). This service was evolved from our push system from the game we worked on and when it just did pushes only and no other logic it could push at 1m/sec in batches, but its responsibility has changed. Context matters :)