How Discord supercharges network disks for extreme low latency
discord.com
discord.com
NVMe - 'non-volatile memory express'. This is a _protocol_ that storage devices like SSDs can use to communicate with the external world. It is orthogonal to whether the disk is attached to your motherboard via pcie, or it lives somewhere else in the datacenter and uses a different transport layer. For example, local NVME SSDs will use PCIe transport. But you can also have NVME-over-TCPIP, or NVME-over-fiber-channel. Or NVME-over-Fabrics. Many of these are able to provide significantly lower latencies than a millisecond.
As a concrete example, AWS `io2-express` has latencies of ~0.25-0.5ms, though i'm not sure which technology it's using.
nvme over fabrics: https://www.techtarget.com/searchstorage/definition/NVMe-ove... many interesting presentations on the official nvme website: https://nvmexpress.org/education/documents-and-videos/presen... aws: https://aws.amazon.com/blogs/storage/achieve-higher-database...
The normal course of action for this is usually to have a raid array over all your nvme disks, but since google just migrates your VM to a machine that has good disks, doing that is useless.
Really this whole article is "we are going to keep using google cloud despite their storage options being unfit for our purpose and here's how".
And that's called hacking. Welcome to Hacker News. Discord's engineers are gods among mortals for squeezing this kind of latency and reliability out of something Google intended for consumers. Yes I know the way you're supposed to do something like this is to use something like Cloud BigTable where you probably have to call a salesperson and pay $20,000 before you're even allowed to try the thing. But Discord just ignored the enterprisey solution and took what they needed from the consumer tools instead. It reminds me of how early Google used to build data centers out of cheap personal computers.
In such a world, perhaps those who remember there's a world of flexibility and power within the OS could be seen as welding some supernatural power.
Hmm, I read it more like they figured out a way to scale their existing storage a bit better while making Google eat the cost.
It’s not like they pay extra if they wear out the disks sooner.
• PCIe Transport specification
• Fibre Channel Transport specification (NVMe-oF)
• RDMA Transport specification
• TCP Transport specification
Separation of storage from transport is a huge game changer. I am really hoping NVMe-over-fibre really takes off. But I'd suggest people would first see that in on-prem deployments before you see it in the cloud-hosting hyperscalers.
More on the 80,000 foot view of what's going on in NVMe world is covered in link below. But there are tech specs you can read over if you are so interested in how exactly it works.
https://nvmexpress.org/nvm-express-announces-the-rearchitect...
[EDIT: Also, the GCP persistent disks are not NVMe-oF as far as we know. They seem to be iSCSI based off Colossus/D. see: https://news.ycombinator.com/item?id=21732387]
The magic here is `--write-behind`/`--write-mostly`[1] in `mdadm`. I mean, that's the only method that I can think of here. This is an old dark magic that prevents (though not entirely) reading from a specific drive.
[1]: https://raid.wiki.kernel.org/index.php/Write-mostly
TBH, in general, I don't think it's a good option for databases. The slow drive does cause issues, for it is literally slow. The whole setup slows down when the write queue is full, reading from the write-mostly device can get be extremely slow thanks to all the pending writes, and the slow drive will wear out quicker thanks to the sustained high load (though this one should not apply in this specific case).
So you mostly don't want to use a single HDD as your mirror. For proper mirroring, you need another array of HDDs, which will end up being a NAS with its own redundancy. That's a large critical pain in the butt, but this is necessary for industrial grade reliability, and also allows making your "slow drive" faster in the future.
In this specific case, it's pretty well played. They get 4 billion messages per day, which is roughly 46k per second. Assuming each write requires updating at least one block - 4kb - the setup needs to sustain at least 185 MB/s, which is clearly beyond a single HDD. Google Persistent Disk is a kind of NAS anyway, so that perfectly aligns with the paragraph above.
It would be a fun exercise to reimplement Discord in AWS...or with FoundationDB.
In any case before revving the hardware I'd want to know how ScyllaDB actually is supposed to perform. I mean, their marketing drivel says this right on the main page: "Provides near-millisecond average latency and predictably low-single-digit P99 response times."
So why are they fucking with disks if ScyllaDB is so good? I mean, back in the day optimizing drive performance was like step 1. Inside tracks are faster, and make sure you don't saturate your SAS drives/Fiber Channel links by mistake. It's fun to do, but you could always get better performance by getting the software to not do dumb stuff. Seriously.
The database layer isn't magic. The database can't give you low-single-digit P99 response times, if a single I/O request can stall for almost 2ms.
That said, I don't think AWS would fare any better here as the infrastructure issue is the same. Networked EBS drives on AWS are not going to be magically faster than networked PD drives on GCP. The bottleneck is the same, the length of the cable between the two hosts.
But up until the last year or two, you couldn't get anywhere near that with EBS and I'm sure as hardware advances, EBS will once again lag and you'll need to come up with similar solutions to remedy this.
Also, I guess AWS would fight them a little less here: the lack of live migrations at least means that a local failed disk is a failed disk and you can keep using the others.
[0]: https://docs.aws.amazon.com/AWSEC2/latest/UserGuide/provisio...
It's really quite amazing to me that HFT reduces RPC latency by about three orders of magnitude, I feel like there are lessons from there that are not being transferred to the tech world.
I've done this in past roles on AWS with their i3, etc. family with local attached storage and didn't use EBS.
> essentially a write-through cache, with GCP's Local SSDs as the cache and Persistent Disks as the storage layer.
'Discord runs most of its hardware in Google Cloud and they provide ready access to “Local SSDs” — NVMe based instance storage, which do have incredibly fast latency profiles. Unfortunately, in our testing, we ran into enough reliability issues'
The ideal here is to buy similarly specced drives from multiple manufacturers to reduce the risk. At the very least buy from multiple suppliers to reduce your risk of getting drives from the same batch if this is something you're going to care about.
https://www.zdnet.com/article/hpe-says-firmware-bug-will-bri... [2020]
What is it about NVME that means they shouldn't fail at a high rate? I don't understand how the protocol should matter much.
Your write load is already split up per-shard in this scenario, so you can horizontally scale out or increase IOPS of the EBS volumes to scale. And you can recover from the hidden secondaries if needed.
And no fancy code written! :)
They say that GCP will kill the whole node, which is probably a good thing if you can be sure it does that quickly and consistently.
If it doesn't (or not fast enough) you'll have a slow node amongst faster ones, creating a big hotspot in your database. Cassandra doesn't work very well if that happens and in early versions I remember some cascading effects when a few nodes had slowdowns.
Edit: Just wanted to add that because they are using Persistent Disks as the source of truth and depending on the network bandwidth it might not be that big of a problem to restore a node to a working state if it's using a quorum for reads and RP >= 3.
Resorting a Node from zero in case of disk failure will always be bad.
They could also have another caching layer on top of the cluster to further mitigated the latency issue until the nodes gets back to health and finishes all the hinted handoffs.
We can either re-build the node by simply wiping its disks, and letting it stream in data from other replicas, or we can re-build by simply re-syncing the pd-ssd to the nvme.
Node failure is a regular occurrence, it isn't a "bad" thing, and something we intend to fully automate. Node should be able to fail and recover without anyone noticing.
But by bad I meant when a Node is based on local disks, that Cassandra and ScyllaDB usually recommends.
Depending on the time between snapshots and the restore from snapshot process (if there are even snapshots...) can be problematic.
Bootstrapping nodes from zero depending on the cluster size (Big data nodes while not recommended are pretty common) could take days in Cassandra because the streaming implements was (maybe still is?) very bad
Thanks.
It's also worth mentioning that in a cluster of N, to recover a node, it simply needs to stream 1/(N - 1) of the data from its neighboring nodes. So when you look at the cluster as a whole, and the strain on each node that is UP and serving traffic, it's insignificant.
uh, you can get a packet from your software on one machine, through a switch and into your software on the second machine in 1-2us if you know what you're doing
(of course, that's without any cloud bullshit like passing the packet through half a dozen levels of virtualisation, then through software defined networks)
therefore, there must at least be several switches between the devices, and realistically, probably one or more routers.
with modern switches it barely matters how many they are in between the machines
This isn’t a pair of servers sitting in a room with a switch in between them on a simple /24 subnet. It’s a gigantic cloud data center with a massive network spanning a huge number of devices.
The simple things you can get to work in a small setup don’t apply at this scale.
it's common to have a dedicated trading network completely parallel and isolated from the main DC network for each upstream connection, of which there are dozens
and these are (at least) double redundant through the entire network, typically using redundant packet arbitration to near guarantee zero drops over a session
Google's network is super-dynamic and flexible, but slow and likely over-contended
ours is the exact opposite
I mean, by definition, a public cloud wouldn’t work without layers of virtualization and SDN.
First, Google has enough money that they can build their entire network out of custom hardware, custom firmware, and patch the kernel + userspace. A datacenter at Google scale is architecturally similar to a supercomputer cluster running on InfiniBand. You will never be able to replicate the performance of Google's network by buying some rackmounted servers from Dell and plumbing them together with Cisco switches.
Second, assuming a reasonably competent design, adding more machines to a network doesn't significantly increase the latency of that network. You'll see better latencies between machines in the same rack than between two racks, but this is a matter of single microseconds rather than milliseconds. Additional latency from intermediate switches is measured in nanoseconds.
Third, Google publishes an SLA on round-trip network latency between customer VMs at <https://cloud.google.com/vpc/docs/vpc>. Their "tail latencies less than 80μs at the 99th percentile" translates to ~40μs for one way, and honestly for customer VMs a lot of that happens in the customer kernel + virtualization layer. A process running on bare-metal, such as a kernel reading a remote network block device, can (IIRC) expect single-microsecond latencies to get one packet onto a nearby machine.
well yes, not true with modern switches that support cut-through forwarding
it's super-common in our space to bypass the kernel entirely, writing into the NIC buffers directly with prepared packet headers, and the card has pushed part of the packet out onto the wire, through switches and into the target machine's NIC buffers before it's even finished being written
typical "SLA"s are 0 packets dropped during a session, where a single drop raises an alert that is then investigated
> You will never be able to replicate the performance of Google's network by buying some rackmounted servers from Dell and plumbing them together with Cisco switches.
and yet, somehow we do quite a bit better (admittedly they are very, very expensive switches)
I get that people that work at Google like to think they're working on problems more advanced than those of mere mortals, but with the latencies you've described we'd be out of business several times over
(not to mention none of the clouds support multicast)
What Google needs is "Big+Cheaper" datacenters, and it has to work with codes written by 100000 different mere morals. What you described is in the "Small+Expensive" field, but with extreme worst-case performance demand.
"Big+Expensive" = Supercomputer "Small+Cheaper" = ??? (Note that the Big+Cheaper solutions not necessarily work for this, as you can't amortize and ignore one-time R&D/ops cost anymore)
The other replies are accepting multi-millisecond latencies as a given, and think that Google's network must be slower than even a basic copper-wired LAN because it's bigger.
My response is something like "just because the network's bigger doesn't mean it's slower".
>> You will never be able to replicate the performance of Google's network
>> by buying some rackmounted servers from Dell and plumbing them together
>> with Cisco switches.
>
> and yet, somehow we do quite a bit better (admittedly they are very,
> very expensive switches)
With respect, if you're in the trading business, your network almost certainly contains custom hardware. I bet it looks a lot closer to Google's than it does to the guy plugging cat5e into a Dell.For reads (unless there is another bottleneck), by 10x'ing the parallelism, an application can compensate for 10x the average latency and still deliver the same throughput.
Scylladb talked about their IO scheduler a good deal:
https://www.scylladb.com/2021/04/06/scyllas-new-io-scheduler...
https://www.scylladb.com/2022/08/03/implementing-a-new-io-sc...
Naively, it seems like it should have been able to compensate for the Persistent Disk latency for read heavy workloads.
Very odd - notice how the latency is high when the read IOPS are low. When the read IOPS climb, 95th percentile latency drops.
Looks like there is a constant rate of high latency requests, and when the read IOPS climb, that constant rate moves to a higher quantile. I'd inspect the raw results but they're quite big: https://github.com/scylladb/diskplorer/blob/master/latency-m...
Background: https://github.com/scylladb/diskplorer/
While I don't know if Google offers fast storage, 0,5ms is slow compared to faster drives. 0,015ms is more around a realistic latency for faster drives with a high queue depth.
> We also lose the ability to create point-in-time snapshots of an entire disk, which is critical for certain workflows at Discord (like some data backups).
No you don't. There are multiple filesystems to choose from that support snapshots for free. Zfs is awesome.
Also, if each database has up to 1TB data and reads are the important part, why not use servers with enough memory to cache the data? That is the normal way to speed up databases that are read heavy.
https://en.wikipedia.org/wiki/Write_through_cache#Writing_po...
With a write-through cache, writes are slow, but the cache is never stale. With write-back, writes can be faster, but the caches can be incoherent.
I have LVM2 set up to use my SSD as a cache for my spinner. I think it's write-through mode. Maybe.
I am curious as to why they can’t tolerate Disk failures though, and why they need to use the Google persistent storage rather than an attached SSD for their workload.
I would have expected a fault tolerant design with multiple copies of the data, so even if a single disk dies, who cares.
Because otherwise, it would make sense to just bounce the node entirely and let it build itself from the replicas.
As a database guy, the number of edge cases inherent in this model is pretty scary. I can think of several scenarios where this is really bad for your system:
* persistent data going back in time (lost the SSD, so you reload from network) * inconsistent view of the data - you wrote two pages in a related transaction. One of them is corrupted, you recover older version from the network.
That is likely to violate a lot of internal assumptions.
https://github.com/mingzhao/dm-cache
https://man7.org/linux/man-pages/man7/lvmcache.7.html
https://www.kernel.org/doc/Documentation/bcache.txt
https://github.com/facebookarchive/flashcache
https://github.com/stec-inc/EnhanceIO
Doing RAID-1 (or RAID-10) style solution certainly is reasonable, but I don't understand what's novel here.
...and IIRC this is the kind of thing that btrfs and zfs were managing explicitly...
How does md handle a synchronous write in a heterogenous mirror? Does it wait for both devices to be written?
I'm also curious how this solution compares to allocating more ram to the servers, and either letting the database software use this for caching, or even creating a ramdisk and putting that in raid1 with the persistent storage. Since the SSDs are being treated as volatile anyways. I assume it would be prohibitively expensive to replicate the entire persistent store into main memory.
I'd also be interested to know how this compares with replacing the entire persistent disk / SSD system with zfs over a few SSDs (which would also allow snapshoting). Of course it is probably a huge feature to be able to have snapshots be integrated into your cloud...
One of the reasons an LRU cache like dm-cache wasn't feasible was because we had a higher than acceptable bad sector read rate which would cause a cache like dm-cache to bubble up a block device error up to the database. The database would then shut itself down when it encountered an disk-level error.
> How does md handle a synchronous write in a heterogenous mirror? Does it wait for both devices to be written? Yes, md waits for both mirrors to be written.
> I'm also curious how this solution compares to allocating more ram to the servers, and either letting the database software use this for caching, or even creating a ramdisk and putting that in raid1 with the persistent storage. Since the SSDs are being treated as volatile anyways. I assume it would be prohibitively expensive to replicate the entire persistent store into main memory. Yeah, we're talking many terabytes.
> I'd also be interested to know how this compares with replacing the entire persistent disk / SSD system with zfs over a few SSDs (which would also allow snapshoting). Of course it is probably a huge feature to be able to have snapshots be integrated into your cloud... Would love if we could've used ZFS, but Scylla requires XFS.
Is it about the ability to dynamically add extra capacity over time or something?
Our estimations for MTTF for our larger clusters would mean there'd be a risk of simultaneous nodes stopping due to bad sector reads. Remediation in that case would basically require cleaning and rewarming the cache, which for large data sets could be on the order of an hour or more, which would mean we'd lose quorum availability during that time.
> to me your solution sounds absolutely brutal for needing a complete copy of the remote disk on the local "cache" disk at all times for the RAID array to operate at all, meaning it will be much harder to quickly recover from other hardware failures.) In Scylla/Cassandra, you need to run full repairs that scan over all of the data. Having an LRU cache doesn't work well with this.
Is the bad sector read rate abnormally high? Are GCE's SSDs particularly error prone? Or is the failure rate typical, but a bad sector read is just incredibly expensive?
I assume you investigated using various RAID levels to make an LRU cache acceptably reliable?
It's also surprising to me that GCE doesn't provide a suitable out of the box storage solution. I thought a major benefit of the cloud is supposed to be not having to worry about things like device failures. I wonder what technical constraints are going on behind the scenes at GCE.
Yes, incredibly error prone. Bad sector reads were observed at an alarming rate over a short period - well beyond what is expected if you were to just buy an enterprise nvme and slap it into a server.
> GCP provides an interesting "guarantee" around the failure of Local SSDs: If any Local SSD fails, the entire server is migrated to a different set of hardware, essentially erasing all Local SSD data for that server.
I wonder how md handles reads during the rebuild, and how long it takes to replicate the persistent store back onto the raid0 mirror.
(I know that live migrations are at least in theory possible, but I don’t know why GCP would go through all the effort)
(I’m also making a lot of assumptions about things I am not an expert in)
When hardware fails, the instance is migrated to another machine and behaves like the power cord was ripped out. It's possible they go down this path for failed disks too, but it's feasible that it is implemented as the disk magically starting to work again but being empty.
You can read more about GCP live migrations here: https://cloud.google.com/compute/docs/instances/live-migrati...
You can read more about GCP live migrations here: https://cloud.google.com/compute/docs/instances/live-migrati...
This seems extremely dangerous as nothing notifies the OS to unmount the filesystem and flush its caches, leading to trashing of the new disk as well. The only way to recover would be to manually unmount, drop all IO caches, then reformat and remount.
That said, from further reading of the GCP docs, it does sound like if they detect a disk failure they will reboot the VM as part of the not-so-live migration.
1. https://access.redhat.com/documentation/en-us/red_hat_enterp...
> while achieving significantly higher throughputs and lower latencies [compared to Cassandra]
Do they really get all that just because it's in C++? Anyone familiar with both of them?
While C++ does have an advantage of raw performance it's ScyllaDB's seastar implementation that helps a lot. Think of every core in the machine as a node so there is no context switching and better use of the cpu cache. More than that the ScyllaDB team are extremely performance focused something that I can't say for Cassandra
My takeaways:
- Cloud vendors should offer a hosted "Superdisk" so users don't have to implement themselves.
- Reading good engineering blog posts can save you the trouble of having to re-learn their experiences!
CPU overhead comes into play when you’re doing parity on a software raid setup (like md or zfs) such as in md raid5 or raid6.
If they needed data scrubbing at a single host level like zfs offers, then probably CPU would be a factor, but I’m assuming they achieve data integrity at a higher level/across hosts, such as in their distributed DB.
Using a RAID manager or a filesystem for this does not seem optimal.
It's certainly not the place to talk if you value your and your peers' safety.
http://widgetsandshit.com/teddziuba/2010/10/taco-bell-progra...
> Here's a concrete example: suppose you have millions of web pages that you want to download and save to disk for later processing. How do you do it? The cool-kids answer is to write a distributed crawler in Clojure and run it on EC2, handing out jobs with a message queue like SQS or ZeroMQ.
> The Taco Bell answer? xargs and wget. In the rare case that you saturate the network connection, add some split and rsync. A "distributed crawler" is really only like 10 lines of shell script.
I generally agree, but it's probably only 10 lines if you assume you never have to deal with any errors.
I will 100% agree that it has disadvantages, but it's unfair to level the above at shell scripts, for most of your complaint, is about poorly coded shell scripts.
An example? sysvinit is a few C programs, and all of it wrapped in bash or sh. It's far more reliable than systemd ever has been, with far better error checking.
Part of this is simplicity. 100 lines of code is better than 10k lines. "Whole scope" on one page can't be underestimated for debugging and comprehension, which also makes error checking easier too.
It’s not about “can bash do it” it’s about “is there a huge ecosystem of tools, which we are probably already using in our organization, that thoroughly cover all these issues”.
If wget fails, you don't retry... at least not until next run.
And wget (or curl, or others) do reply with return codes which indicate what kind of error happened. You can also parse stderr.
Of course you could programmatically handle backoff in bash too, but.. why? Wget is very good at that. Very good.
===
In terms of 'junior dev', a junior dev can't contribute to much without ramp up first. I think you mean here, 'ramp up on bash' and that's fair... but, the same can be said for any language you use. I've seen python code with no error checking, and a gross misunderstanding of what to code for, just as with bash.
Yet like I said, I 100% agree there are issues in some cases. What you're saying is not entirely wrong. However, what you're looking for, I think, is not required much of the time, as wget + bash is "good enough" more often than you'd think.
So I think our disagreement here is, how often your route is required.
And while it sounds Unixy to let wget do its thing, a fully baked program like that is much less “do one thing and do it well” than the http utilities in general purpose programming languages.
Is grey hair a new dependency on using wget? Upstream hasn’t updated their dependencies apparently.
Plenty of stuff runs in production as shell scripts at Google - if you’re a dev, you just don’t notice it often.
Most of the rest of the glue is Python.
These tools were processing and generating phone bills in the 1980s with computers with less computing power than your watch.
Then run your little pipeline in a loop until it stops making progress (`find | wc` doesn’t increase.) Either it finished, or everything that’s left as input represents one or more classes of errors. Debug them, and then start it looping again :)
The issue here is that your code has no real-time adaptability. Many backends will scale with load up to a point then start returning "make fewer requests". Normally, you implement some internal logic such as randomized exponential backoff retries (amazingly, this is a remarkably effective way to automatically find the saturation point of the cluster), although I have also seen some large clients that coordinate their fetches centrally using tokens.
You know how you can rate-limit your requests? A forward proxy daemon that rate-limits upstream connections by holding them open but not serving them until the timeout has elapsed. (I.e. Nginx with five lines of config.) As long as your fetcher has a concurrency limit, stalling some of those connections will lead to decreased attempted throughput.
(This isn’t just for scripting, either; it’s also a near-optimal way to implement global per-domain upstream-API rate-limiting in a production system that has multiple shared-nothing backends. It’s Istio/Envoy “in the small.”)
Having built several large distributed computing systems, I've found that the inner client always needs to have a fair amount of intelligence when talking to the server. That means responding to errors in a way that doesn't lead to thundering herds. The nice thing about this is that, like modern TCP, it auto-tunes to the capacity of the system, while also handling outages well.
> Having built several large distributed computing systems, I've found that the inner client always needs to have a fair amount of intelligence when talking to the server.
I disagree. The goal should be to make the server behave in such a way that a client using entirely-default semantics for the protocol it’s speaking, is nudged and/or coerced and/or tricked into doing the right thing. (E.g. like I said, not returning a 429 right away, but instead, making the client block when the server must block.) This localizes the responsibility for “knowing how the semantics of default {HTTP, gRPC, MQPP, RTP, …} map into the pragmatics of your particular finicky upstream” into one reusable black-box abstraction layer.
We Unix graybeards may be used to xargs, grep and wget. The next generation of developers are learning how to construct pipelines from step functions, sqs, lambda and s3 instead. And coming as someone who really enjoys Unix tooling, the systems designed with these new paradigms will be more scalable, observable and maintainable than the shell scripts of yore.
I think cloud gets much maligned — but all the serious discussions with, eg, AWS employees work from this paradigm:
- AWS is a “global computer” which you lease slices of
- there is access to the raw computer (EC2, networking tools, etc)
- there are basic constructs on top of that (SQS, Lambda, CloudWatch, etc)
- there are language wrappers to allocate those for your services (CDK, Pulumi, etc)
…and you end up with something that looks surprisingly like a “program” which runs on that “global computer”.
I know that it wasn’t always like that — plenty of sharp edges when I first used it in 2014. But we came to that paradigm precisely because people asked “how can we apply what we already know?” About mainframes. About Erlang. About operating systems.
I think it’s important to know the Unix tools, but I also think that early cloud phase has blinded a lot of people to what the cloud is now.
All the crawler needs to be is a quick crawler script, a Typescript definition of resources, and you get all the AWS benefits in two files.
Maybe not “ten lines of Bash” easy, but we’re talking “thirty lines total, with logging, retries, persistence, etc”.
In their case, I'm wondering why the host failure isn't handled at a higher level already. A node failure causing all data to be lost on that host should be handled gracefully through replication and another replica brought up transparently.
In any case, their usage of local storage as a write through cache though md is pretty interesting, I wonder if it would work the other way around for reading.
[0] https://discord.com/blog/why-discord-is-switching-from-go-to...
[1] https://discord.com/blog/using-rust-to-scale-elixir-for-11-m...
> Rewrite it in Rust
Idk, I don't see the irony here
Related, I hope GCP is listening and builds an "out-of-the-box" solution that automatically combines this write-through caching solution into one offering. Just like I shouldn't have to worry (much) about how the different levels of RAM caching work on a server, I shouldn't have to worry much about different caching layers of disk in the cloud.
I wish companies would stop inflating their numbers by citing "per day" statistics. 4 billion messages per day is less than 50k/second; that sort of transaction volume is well within the capabilities of pgsql running on midrange hardware.
It can in theory work, but the real world would make this the most unstable platform of all the messaging platforms.
Just one vacuum would bring this system down, even if it wasn't an exclusive lock... Also I would be curious how you would implement similar functionality to the URL deep linking and image posting / hosting.
Mind you the answers to these will probably increase average message size. Which means more write bandwidth.
Some bar napkin math shows this would be around 180GiB per hour, 24/7, 4.3TiB per day.
Unless all messages disappeared within 20 days you would exceed pretty much any reasonable single-server NVME setup. Also have fun with trim and optimizing NVME write performance. Which is also going to diminish as all the drives fail due to write wear...
That's my point. Per-second numbers are far more useful than per-day numbers.
Useful for what exactly? You're basically accusing the company of inflating their numbers (??) as if they are deliberating lying about their scale.
1) avg mean (for throughput estimates), 2) p50 (typical customer / typical load), 3) p90 / p99 / p99.99 (whatever you think your tail is) 4) p100 (max, always useful to see and know, maybe p0).
Or throw in a histogram or kernel density estimate, sometimes there are really interesting patterns.
I’ve never seen a technical blog give such traffic or latency details though.
Edit: reading other comments, please do not read this as diminishing this blog post or Discord, great clear writing and impressive solution to an interesting problem.
It's a blog...it IS marketing.
>Our databases were serving around 2 million requests per second (in this screenshot.)
I doubt pgsql will have fun on mid-range hardware with 50k writes/sec, ~2 million reads/sec and 4 billion additional rows per day with few deletions.
> 50k/second
Yes, 50k/second for every minute of the day 365/24/7. Very few companies can quote that.
Not to mention:
- Has complex threaded messages
- Geo-redundancies
- Those message are real-time
- Global user base
- Unknown told of features related to messaging (bot recations, ACL, permissions, privacy, formatting, reactions, etc.)
- No/limited downtime, live updates
Discord is technically impressive, not sure why you felt you had to diminish that.
I don't see anything on their messaging specifically, just assuming they would have something similar.
https://discord.com/blog/how-discord-handles-two-and-half-mi...
Data storage is going to be multi-regional soon, but that's just from a redundancy/"data is safe in case of us-east1 failure" scenario -- we're not yet going to be actively serving live user traffic from outside of us-east1.
1. Cherry pick a piece of info 2. Claim it's not that hard/impressive/large 3. Claim it can be done much simpler with <insert database/language/hardware>
It's not even the most interesting metric about our systems anyway. If we're really going to look at the tech, the inflation of those metrics to deliver the service is where the work generally is in the system --
* 50k+ QPS (average) for new message inserts * 500k+ QPS when you factor in deletes, updates, etc * 3M+ QPS looking at db reads * 30M+ QPS looking at the gateway websockets (fanout of things happening to online users)
But I hear you, we're conflating some marketing metrics with technical metrics, we'll take that feedback for next time.
That said, I don’t totally disagree with you that maybe collocating would be worth it. But they were very clear that they like offloading the work of durable storage to GCP. And they found a pretty elegant solution to achieve the best of both worlds.
Also, FWIW, Discord does offload some workloads to their own hardware. (Or so I’m told, I don’t work there but I know people who do)
It is the height of chutzpah to talk about supercharging and extreme low latency, when the end user experience is so molasses slow.
> Local SSDs in GCP are exactly 375GB in size
Ouch! My SSD from a decade ago is larger than that, and probably, more performant.
* https://www.unrealircd.org/docs/Channel_history
* https://docs.inspircd.org/3/modules/chanhistory/
* https://github.com/ergochat/ergo/blob/master/docs/MANUAL.md#...
And they are not niche servers; UnrealIRCd and InspIRCd are the top two in terms of deployments according to https://www.ircstats.org/servers
If it's a sequential write (by downloading the entire blockchain), you will still be bottlenecked by the throughput of the underlying disk.
If it's sequential reads (in between writes), the reads can be handled by the cache if the location is local enough to the previous write operation that it hasn't been evicted yet.
If it's random unpredictable reads, it's unlikely a cache will help unless the cache is big enough to fit the entire working dataset (otherwise you'll get a terrible cache hit rate as most of what you need would've been evicted by then) but then you're back at your original problem of needing a huge SSD.
Unless you're running a full archive node, you don't need 2tb. My geth dir is under 700gb. Do a snap sync then turn on archive mode if you only need data moving forward from some point in time.
In Azure, I don't have to worry about any of this. There's a flag to enable read-only or read-write caching for remote disks, and it takes care of the local SSD cache tier for me! Because this is host-level, it can be enabled for system disks too.
Disclosure: I work at ScyllaDB.