Scaling Datastores at Slack with Vitess
slack.engineering
slack.engineering
It seems like a well executed combination of Vitess/Citus + read replicas pointing to decoupled page storage would result in a fully featured non-proprietary database that scales all the way down to your laptop and all the way up to the heaviest OLTP use cases. If that existed, only truly masochistic people would continue using things like DynamoDB.
Which is probably why AWS doesn't support it, they make less money and have less lock in on their clients.
> Each keyspace is a logical collection of data that roughly scales by the same factor — number of users, teams, and channels. Say goodbye to only sharding by team, and to team hot-spots!
Hi Sugu!
I wouldn’t say that killed my first startup, but it certainly didn’t help.
Database like CockroachDB and TiDB help here.
You can easily buy a machine with several TB of RAM, several hundreds of TB of SSDs in RAID giving you millions of IOPS, quad-socket 256 cores. How likely is it that a single machine cannot handle a single customer?
You either need an enormous client base or are doing something very specialized that requires tons of memory or cpu per customer.
I would think 99% of all internet companies could likely handle everything on a single massive machine, or two, for redundancy. Assuming they were somewhat optimized.
As your database gets into tens or hundreds of TBs, now backups and restores take hours or days to complete. Any DR scenario becomes an existential threat.
For large tables, schema changes can take weeks. This often results in developers doing their own custom table sharding strategies, even though they might still be on a single machine.
I'm not saying that you can't architect around these problems, but operating a db at massive scale on a single machine isn't the obvious win that it seems to be.
Damn, even for a service as popular as Slack, that’s significantly more than I expected. Slack has ~12-13 mil DAUs, right? I assume at peak time of day, maybe 2-3 million actively using the product at the same time? If that’s a fair assumption, that’s roughly 1 MySQL query per second, per active customer - seems like a fair bit? I wonder if they do polling (instead of websockets)?
At my work we have roughly 1 order of magnitude fewer DAUs, but roughly 2 orders of magnitude fewer QPS. And we also have chat areas of the product.
As of September 2019, Slack reported 12 million DAUs [0]. Of course, that is pre-pandemic, and the only figures that Slack has provided regarding demand in 2020 are these tweets [1] from Slack CEO Stewart Butterfield back in March. In those tweets, it mentions that Slack was serving 12.5M simultaneously connected users.
[0] https://slack.com/blog/news/work-is-fueled-by-true-engagemen...
While that's like 1/7th of Slack's traffic, it is served by a single PostgreSQL cluster.
That's the comment.
Slack are getting 300K WRITES per second, though. Which does mean you have really no choice but to shard - read replicas are nice for scaling reads, but do nothing for writes, and that’s A LOT of writes.
Right now the cluster is only 8 instances, and the write-only traffic on the master spikes to 70K QPS.
So I agree with you --pretty beefy hardware, and it'd probably require sharding before reaching Slack's scale, but still quite impressive.
Vitess plans to support PostgreSQL, but there is no work on that at this time.
If all you want is HA and FO there are simpler solutions like pgpool and pgbouncer etc.
> stolon - PostgreSQL cloud native High Availability https://github.com/sorintlab/stolon
- crunchydata - the thing crunchy is based on (well actually they are more based on the underlying things like patroni): zalando postgres operator https://github.com/zalando/postgres-operator, really really solid stuff