Scaling Databases at Activision [pdf]
static.sched.com
static.sched.com
That's like 10 tall servers for their peak QPS? (500k QPS on slide 19)
Lots of natural sharding points here too: company, game, usage of data, etc.
This is assuming you avoid or go light on expensive features like foreign keys (as Vitess already does).
The "scale magic" in Spanner (Big Table) inspired DB's are just hidden automation of traditional sharding-
CockroachDB: Sharded indexes, then runs map reduce for you.
Scylla/Cassanda: Sharded indexes again, but more limitations for speedups: Eventually consistent. You don't have fast delete- only update (discord uses tombstones). JOIN's are "in app" only.
Vitess: Proxy that dismantles/routes your query to the correct server. This scales, but is eventually consistent. JOIN's on co-located data, or "in app".... in simple terms, Vitess is like an externally managed "in app" query parser/router.
Data is divided into chunks and the routers route queries based on where your query lands on chunk ranges.
In my experience me this isn't close to being possible with a single Postgres box. There is no way around sharding this workload.
Vitess has different consistency models, it's not as eventually consistent as you describe.
I'm not quite sure if this is a common problem or if it's our workload. It always made me wonder if moving pooling into Postgres is really a good idea.
Golang and ASGI Python (FastAPI, Starlette, Sanic) can do a small connection pool per worker, each worker can handle thousands of users at once.
PHP in 2023 has database connection persistence using FPM + PDO (effectively, one big connection pool shared between all workers on a server). Again, eliminates connection overhead- works well.
Just remember to raise your max_connections in postgres, like in this: https://gist.github.com/gnat/cfe3754c3dc817c7fb8b2225ef4db62...
That’s a pretty big one. If you know you need the scale, the trade off might be worth it. But I’m 99% of projects, Postgres is more than adequate and has great flexibility.
That’s why people use these systems and why so many very smart people spend so much time trying to make them better
And regarding correctness and iteration speed.. yes sure but it's easy to say that you are correct when you only run on a single node. The reason why those distributed solutions are less correct is because it's harder to be correct. For example in vitess you end up having to do a bunch of manual work with their sequences for every table. In scylladb afaik you just trust that the timeuuids are unique or you do the same and pregenerate bigints into a table.
My point is: To say that mysql/postgres are safer because they are correct and you can iterate faster is only true if you go down a very narrow path with your orm and you stay on a single node. You are essentially postponing all the pain. At some point you'll stop using the orm and write your own, at some point you'll stop using the single node.
Scylla is great for specific workloads, but it’s a niche tool that makes sense for specific high volume workloads but I wouldn’t recommend it as a general purpose application development database. It will slow you down a lot and a lot of really simple features will be hard to implement.
Also I’m not talking about ORMs. Writing SQL by hand in your application is 10-100x more productive for feature development than using something like Scyllas APIs directly IMO.
SQL is for relational databases. That's a different paradigm to CQL. The similarity is superficial. CQL is not a dialect of SQL. Not even close.
Of course, sometimes they do have actual performance issues, but 99 out of 100 times they just don’t know what an index is.
The main issue I've run into is deployment. It's just very rough compared to CockroachDB or Postgres, which are single binary or very close to it.
If ScyllaDB became as easy to set up as CockroachDB (or close-ish), they'd have a lot more users, I think.
You can make Postgres setup this easy: https://gist.github.com/gnat/cfe3754c3dc817c7fb8b2225ef4db62...
CockroachDB, also dead easy: https://www.cockroachlabs.com/docs/v22.2/start-a-local-clust...
ScyllaDB is an enormous mess of daemons and ports: https://university.scylladb.com/courses/scylla-operations/le...
As for the "huge chinese" databases, not to sound like a conspiracy, but CockroachDB has outperformed them in my server load testing (stress test with very basic "SELECT FROM users" and "INSERT INTO users" on KC3000 + Ryzen 7950x)- so I'm beginning to seriously question all of the self-published benchmarks from mainland china.
That, and when one considers they have a similarly rough deployment story compared to ScyllaDB, is it maybe worth just going directly to ScyllaDB in the first place?
One other damming thing: I've also observed that all Spanner (Big Table) inspired databases have lower throughput compared to Postgres when under 5-ish servers: Aka you only begin to see throughput benefits after you have a sizeable cluster.
ScyllaDB / Cassandra is a different architecture, though.
Heck even Vitess, pretty rough deploy story as well (they really push k8s on you), but that's why PlanetScale exists, IMO.
But you are right about how important it is to have a simple deployment, it's true for all of the databases. Sure there is docker but that usually has its own problems, like with scylla where it wouldn't properly forward ports for whatever reason so I could never connect to it. Then I had to do all kinds of manual fixes and install it on a VM to edit configs to make it run. This is really the stuff that makes you want to use crdb and be done with it all, even though it's slower and eventually requires payment.
Write Postgres in a way that's fully compatible with CockroachDB, with the intention of transitioning to CockroachDB if needed.
Basically: 1. BIGINT or UUID or TEXT or composite Primary Keys only. 2. Limited use of sequential indexes. (Ex: timestamp. Not the end if the world but creates hotspot, will require a hash sharded index) 3. Avoid advanced features such as stored procedures. 4. Avoid or very light use of foreign keys. (Performance)
Going straight to CockroachDB seems logical as well, but you'll hit the performance wall much sooner on 1-2 servers, and will require a cluster to match it.
but scylladb really has so many weird footguns. I just discovered that you can't even query for NULL values. It just doesn't work, they have no support for WHERE x = NULL even with an index. That's crazy to me, if I ever make a mistake and have a few NULL'ed values then I can't even find them again so I can set them to an empty string to make them queryable. I just don't know, I almost feel like scylladb has too many tradeoffs.
For the record, CockroachDB does not store null either, but it allows querying null.
A workaround may be stepping through the full table manually using the PK (which can never be null) checking for null "in app" and cleaning them that way.
Noticed INSERTS are treated as UPSERT unless you use IF NOT EXISTS... ScyllaDB doesn't give a f__k if it's overwriting an existing row, lol. I can live with the UPSERT footgun, though. Agreed the null footgun is far more annoying.
SELECT * FROM users;
Then you can plug your PK into.
UPDATE users SET address='' WHERE user_id IN (9affeac3-5d92-111d-779c-55eb6d78a806, ...) IF address=null;
Then you can do your SELECT to do your repair.
SELECT * FROM users WHERE address='';
_______________________
Another footgun: No default values for columns in CREATE TABLE. Very annoying.
Another footgun: Only "=" and "IN (...)" is supported in WHERE for partition keys. https://docs.scylladb.com/stable/cql/dml.html#the-where-clau...
This means, you cannot do UPDATE WHERE user_id > 0 ....
I'm starting to really appreciate the effort CockroachDB went to, to make their version of "CQL" postgres compatible, even though architecturally they are both key-value store databases under the hood... it's just a shame CockroachDB just has far lower performance because of enforced consistency and no way to turn it off.
To be fair, ScyllaDB is removing footguns with every new version, just not fast enough for my taste.
In contrast to CockroachDB, Postgres, etc, where Consistency in CAP is always enforced and there's no way to turn it off.
You turn it on/off by simply specifying IF NOT EXISTS which is insanely simple: https://www.scylladb.com/2020/07/15/getting-the-most-out-of-... (called "lightweight transactions" or LWT)
I wish CockroachDB had a way to trade consistency for speed when desired.
Another sidenote, in my limited investigation so far- when using LWT in ScyllaDB: performance of CockroachDB INSERT and ScyllaDB INSERT both line up fairly evenly. Still investigating, but this makes sense- Only when you avoid IF NOT EXISTS, ScyllaDB pulls away massively in performance.
> Citus also exists for Postgres but their docs basically tell you that it's only recommended for analytics
Could you point me at the docs that made you think that? We find Citus very good at multi-tenant SaaS apps (OLTP), IoT (HTAP) workloads and analytics (OLAP).
> And it sounds like all it's doing is a basic master-slave postgres setup with quite a few manual things you have to do to even benefit (manually altering tables to make them sharded/partitioned)
This is true for now. We are looking into ways to make onboarding easier. That said, the time spent on defining a good sharding model for your data often leads to very good perf characteristics. Regarding the architecture, I personally find Citus closer in spirit/design to what Vitess is doing. Additionally, every node in the cluster is able to take both writes and reads, so I don't see the parallel to a basic primary/secondary setup.
There is a middle ground, storage systems/databases that both scale and provide strong consistency. FoundationDB, Citus, etc.
SQL databases are not only easier to query, but also enforce data consistency. While querying is their core job, it's hardly the only one.
Which is a huge downside.
> But so what, at least it's fast
So is an RDBMS for most read/write workloads. Few will see an overall benefit from Cassandra and the labor needed to replicate the conveniences offered by an RDBMS.
> it's about storing data
It's not about storing data, it's about using data to accomplish business objectives with the promises that entails. You opt for something like Cassandra/Scylla because it solves a technical challenge otherwise untenable or more expensive, but it's a poor choice if your problem is reasonably solved with what Postgres can offer.
* Low volume, complicated queries
* High volume, simple queries
In the rare case that I've seen high volume, complicated queries it generally involves a lot of pre-indexing, denormalization, and exotic data structures.
Postgres really shines in the first case, which is probably the dominant case for "data important enough to pay people to work on".
In theory a distributed SQL database can be a huge developer win because of this. Of course there are usually other tradeoffs to consider as well.
It is amazing what Vitess / PS will help you accomplish but they don't talk a lot about tradeoffs you make to get there (which are similar to what you face with sharding MySQL without additional tooling).
I stumbled upon your DB service several times, and considered it several times suggesting it as the DB of choice for my customers (uxwizz.com), but the pricing is a bit confusing/hard to estimate in my case, which makes it hard for me to recommend it. Any suggestions?
On delete cascade, depending on how many rows it cascades to, can be problematic because it's a very long running blocking operation. That's something one might want to do as a background operation and in batches. Although that won't make it faster.
Personally, I find delete on cascade dangerous. I mean, lots of fun for a pen tester, sure…
Intriguing. Can you tell me more?
Sure, I meant an index in which the key is primary. But, on second thought, that's probably me misreading the GP's message.
> On delete cascade, depending on how many rows it cascades to, can be problematic because it's a very long running blocking operation. That's something one might want to do as a background operation and in batches. Although that won't make it faster.
That can definitely be a problem (just like destructor deallocation stampedes in C++ or Rust). Still less risky than cascading manually and asynchronously, I suspect.
You're conflating the concept of a normalized database with insanely slow DB-enforced referential integrity/foreign keys.
Toy example? Sure. Take a well-formed 3NF schema and disable foreign key constraints.
> Take a well-formed 3NF schema and disable foreign key constraints.
I'm familiar with 3NF, but can you expand on how 3NF enables you to remove foreign keys? Or feel free to point me to an article/blog, I don't want to waste your time if it's too much to explain. I did some googling but wasn't sure where to proceed from your post.
user = {user_id, email}
order = {order_id, user_id}
order.user_id is a foreign key to user.user_id. That's a perfectly valid and reasonable way to organize things.Enabling RDBMS-enforced foreign key constraints is the issue. It slows everything down dramatically.
How slower is it really? MySQL automically creates indexes for foreign keys so I don't think it slows down "dramatically", just an additional indexed retrieval?
Do you know of any benchmarks which show the difference?
Something like SQL Server can enforce foreign key constraints if you explicitly tell it your relationships between tables. The downside is that having this referential integrity costs you performance as the database has to check your relations when inserting/updating/deleting rows. E.g. checking that a foreign key is pointing at a valid primary key, checking that you aren't leaving invalid foreign keys when deleting a primary key, etc. This is to prevent you inserting bad data into the database.
You can delete these constraints and still have the exact same behavior so long as your code is correct. It just means that the database isn't going to stop you writing bad data.
My sense is that many typical CRUD apps aren't writing gargantuan volumes of data or making very complex edits, and if they do, it's ok if it takes a second longer. Usually read speed is more of a bottleneck for user-facing applications, but I'm sure there are probably some examples where this tradeoff is worth it.
This is a pretty tough definition of correct, though. Without foreign key constraints you'll have a really tough time dealing with concurrency artifacts without raising your isolation levels, which generally brings larger performance concerns.
My experience is that if you have a moderate amount of foreign keys, a lot of DBMS (not Postgres) will refuse the `ON DELETE CASCADE` (in the diamond case), and you have to do it "manually" anyway (from your query builder).
I’m not talking about using cascade - this applies perfectly well to use of on delete restrict. FKs are more or less the only standard way to reliably keep relationships between tables correct without raising up the isolation level (at least, in most dbs) or doing explicit locking schemes that would be slower than the implicit locking that foreign keys perform.
Table Resources: resourceid, userid, etc
If I want to restrict deletion of a user to only be possible after all the resources are deleted, I'm forced into using higher-than-default isolation levels in most DBs. This has significant performance implications. It's also much easier to make a mistake - for example, if when creating a resource I check that the user exists prior to starting the transaction, then start the tran, then do the work, it will allow insertion of data into a nonexistent user.
select id from users where id = ? for update;
if row_count() < 1 then raise 'no user' end if;
insert into sub_resource (owner, thing) values (?, ?);
commit;
??
If we take postgres as an example, performing the select takes exactly zero row level locks, and makes no guarantees at all about selected data remaining the same after you’ve read it.
edit: my mistake - I missed that the select is for update. Yes, this will take explicit locks and thus protect you from the deletion, but is slower/worse than just using foreign keys, so it won't fundamentally help you.
further edit: let's take an example even in a higher isolation level (repeatable read):
-- setup
postgres=# create table user_table(user_id int);
CREATE TABLE
postgres=# create table resources_table(resource_id int, user_id int);
CREATE TABLE
postgres=# insert into user_table values(1);
INSERT 0 1
Tran 1:
postgres=# BEGIN TRANSACTION ISOLATION LEVEL REPEATABLE READ;
BEGIN
postgres=# select * from user_table where user_id = 1;
user_id
---------
1
(1 row)
Tran 2:
postgres=# BEGIN TRANSACTION ISOLATION LEVEL REPEATABLE READ;
BEGIN
postgres=# select * from resources_table where user_id = 1;
resource_id | user_id
-------------+---------
(0 rows)
postgres=# delete from user_table where user_id = 1;
DELETE 1
postgres=# commit;
COMMIT
Tran 1:
postgres=# insert into resources_table values (1,1);
INSERT 0 1
postgres=# commit;
COMMIT
Data at the end:
postgres=# select * from resources_table;
resource_id | user_id
-------------+---------
1 | 1
(1 row)
postgres=# select * from user_table;
user_id
---------
(0 rows)
You can fix this by using SERIALIZABLE, which will error out in this case.This stuff is harder than people think, and correctly indexed foreign keys really aren't a performance issue for the vast majority of applications. I strongly recommend just using them until you have a good reason not to.
I think it's good practice to enforce consistency rules both in the DB and in code. If you make a mistake in your code, the DB won't allow it, and vice versa.
And as you point out, there are exceptions, like financial data. But not marketing funnels where you might throw everything away.
You should probably have different types of engineers working on such different projects as well.
I’ve seen several teams regret making this assumption. Unless you enforce draconian access control over a database, you’re going to discover that customer support and accounting and biz dev have been quietly relying on it.
Plus, relational databases don't just sit under a single application. There's usually multiple applications/services talking to them. Worse, humans connect to them an do all sorts of things they shouldn't do. That's the whole point of managing referential integrity in the DBMS, since you can only control "just write good code" across so many application domains.
Of course whether the performance tradeoff is worth it is a complicated decision for many of the reasons people have mentioned. But in 20 years of working with relational databases at big companies, I've seen few examples where the performance win exceeded the business risk.
I'm also going to disagree with you a little, again on the issue of reading, if you have foreign keys then you can do optimisations like cutting out chunks of joins. Adding constraints means you know more which you can feed to the optimiser which means often you can get better performance (read performance, that is).
As to why you can't really have foreign keys, basically, sharding schemes like this sit at an intermediate layer between the application and the various DB clusters that are the shards. You CAN have strong data consistency (and useful foreign keys) within a shard because it's all inside a single database; however across shards, you CAN NOT. The sharding layer doesn't perform checks for you, so if you have two logical tables that are partitioned differently across the shards, you can't have a foreign key that will be enforced correctly as the foreign key of a given row in one table may live on a different shard in the other table. The local database within the shard would reject the insert. Transactions across shards can also be tricky to impossible.
For example, a multi-tenant application could have a tenants table with PK id, and an orders table with PK (tenant_id, id), and a FK from orders(tenant_id) to tenants(id). Then CockroachDB would keep the orders with tenant_id = 123 in the same machine where the tenant with id = 123 was.
I have the impression they deprecated interleaved tables, but I can't find more info about it.
https://github.com/cockroachdb/cockroach/issues/52009
Not only deprecated, but completely removed a few years ago as seen in:
Now that Vitess has native schema change tools it's more reasonable to revisit user-friendly, out-of-the-box foreign key support.
> ● SQL query compatibility
> ● Minimal changes to the application
> ● Runs MySQL in backend
> ● Kubernetes native
> ● Provides Kubernetes operator
> We evaluated multiple candidates and chose Vitess
https://github.com/vitessio/vitessBut it's very light on numbers and doesn't show trade-offs or anything.
Would love to hear more about this implementation
Just trying to do basic scaling or other simple tasks sounded like an utter nightmare on their old system. Having autonomic computing at their back seems like an obvious win. Tell it your desired state & let the controller do the job.
The new paradigm for computer operations is so wonderful.
To me this is validation that going the SQL route in almost 99% of new apps is the right way to go. It will be a rare case that you won't be able to scale out given how mature some of these technologies are.
There are very few things that structured NoSql can't do if you ignore reporting. Once you scale out traditional SQL, you aren't going to be using it for reporting either way though.
Automating master failover on MySQL without a human in the loop - in the topology GitHub used back then - is risky.
0: https://github.blog/2018-10-30-oct21-post-incident-analysis/
* How did you tune Vitess for resiliency? What were the tradeoffs and how was the performance?
* How did you migrate from shared in app config to single DB endpoint? (this is an issue we are facing right now)
* What do you mean by some queries being too shard aware? How did you optimize queries for efficient routing with Vitess?
At 30 TBs the hot data set is huge an there is a long tail of queries not loaded into memory.
Plus they were saying 10k+ connections per server. So we are probably looking st 500k+ connections.
There is also a regular influx on unoptimized / badly written queries as people add features/games etc. Plus you can end up with random spikes.
This setup is likely optimized for high read availability. The slides mention the mental overhead of failovers. This is a thing. Try to explain and help each team understand how to do those safely. Downtime are also not an option.
So there is quite some information and context (likely) missing. I used to manage a similar scale zoo of MySQL (similar peak QPS, less total data iirc). I am not sure I would subscribe so much on the Vitess as solution side though. Galera and Group Replication can give you hassle free failovers, add some read replicas and you have a really nice DB setup that can scale horizontally amd vertically until you hit the write limitations per server.
But I think the industry are being pushed hard and dont have much time to produce content for our reading interest.
I would have liked a bit more specifics on their manual sharing solution, and how the vitess implementation differs.
Those queues are usually the result of intentional capacity planning. Overprovisioning is expensive.
I'd imagine they could spin up new servers on demand.
I do feel safe in saying though that the login queues are _definitely_ not hype.
They're designed to constrain a quite-complex distributed system to a login rate that has been load-tested thoroughly, e.g. they know it'll work at that rate.
I always found disabling FK checks weird, as in, once you disable them, isn't the integrity permanently lost? Or does it check again the integrity once it's reactivated?
So it's a story of people who used Kubernetes and MySQL and continued to use Kubernetes and MySQL. The end.
Thanks for the list (and the graph), very interesting!
We had no way to have that kind of confidence about migrating off MySQL to an entirely different database (like TiDB). If we'd been addressing a new storage use case from scratch instead of migrating an existing one, the case for using something other than MySQL might have been stronger.