Benchmark results: Cassandra 4.0, 3.11, Scylla 4.4 [video]
youtube.com
youtube.com
I suspect this talk is a follow up to this blog post, which may be preferable for many.
They are using CitusDB but it's still RDBMS and support most PostgreSQL SQL query without the need to rewrite your query or change your schema.
I dont know if they are still using Vitess or migrated to Google other relational Database "Spanner"
In general I've noticed that people outside Google assume that Spanner is much more popular than it really is.
if you work at google and its one of the option to choose from! Is there any reason to "not" choose spanner for new project and instead use something else?
Any of the above makes your data accessible via standard SQL. You can hire and utilize proficient data analysts, not people writing esoteric queries in whatever flavor of NoSQL, while actually solving business needs.
How are Scylla / Cassandra better?
Meanwhile I'm over here using Scylla and 3 nodes can give me roughly 3.5M TPS for data extraction.
RDBMS are great for certain use cases but they do not scale and if you think they do you don't really know what scale means.
Sure you can cite things like Github using MySql, but i'll also point out they've invested thousands upon thousands of man hours customizing layers on top of MySql to support their sharding which also, oops, had a bug that gave them almost half a day of downtime at one point. This is not a good use of someone's time, GH definitely made the wrong decision. Don't steer other people into making similar mistakes especially given it seems you have a lack of experience on designing systems at scale.
Clickhouse vastly outperforms and outscales Scylla with a fraction of the hardware and cost.
Should I have simply said "SQL databases where data is represented in tables with a defined schema" to simplify the discussion and prevent your ignorant diatribe?
Have you used any of these technologies? Cassandra/Scylla have a defined schema. Did you really not know that?
If you're using Cassandra and you're also not a fool, it's probably because you actually need a lot of scale, so you'll generally do the denormalization up front and you will denormalize everything. If you're writing an event, you will write it several places. For example, suppose you're WalMart recording sales. You might write to: store transactions by store/year/day/hour (the "master" record insofar as you have one), user transactions by user/year/month, product purchases by manufacturer/year/day/hour... When you write the transaction, you write to all of these locations.
Each of these "by X" keys is a shard. Each can be located on a different set of ~3 machines (the number is configurable). Querying involves getting a copy of the ring topology, computing which integer shard-ID the key maps to, figuring out which machines in the ring own that integer, and then asking the machine for a whole bucketful of data, which should be a superset of what you're actually looking for. For something like a user's transactions you'll want to have basically everything there at once, so loading the "order history" page for the past month might be a single query that just returns a report: no joins at all, very fast, super scalable. Other lookup strategies might ask for a range within that bucket (the data within the bucket can be ordered by a single key; often this is a timestamp or time-based UUID). Anything that isn't a simple query of a few buckets like this has to be a map-reduce job and will be slow.
All of this is pain. You should generally not invite pain into your organization. However, if pain has already found you, something like this may be the least painful option.
Not really a superset, cassandra (and scylla, and bigtable, all of which are basically copying bigtable's model) each try very hard not to read any extra data at all, and can often return approximately the exact data requested, modulo serialization data which is usually fitting in ~compression chunk size (64k) + checksum.
> If you're using Cassandra and you're also not a fool, it's probably because you actually need a lot of scale
Cassandra also gives you very literally the most control over CAP tradeoffs of any database in the industry.
If you have 100 machines per DC in each of 10 dcs, what happens when one machine is offline? one rack? one dc? 2 dcs separated from 8? 6 dcs separated from 4? There's no single answer in cassandra (depends on replication factor, consistency of writes, consistency of reads, all of which are tunable, with 2 of those being tunable PER QUERY), the CAP tradeoffs are yours and yours alone. That flexibility is powerful for power users (it's also confusing for novices, which is unfortunate).
But to your first point, yes, the point is scale. The lack of opinions and deliberate functionality are designed to enable it to scale to thousands of hosts, potentially petabytes of data, trivially accessible in a single SQL-like CQL query, with realistic read latency < 1ms mean/avg and < 5ms p99 for a tuned workload where you know what you're doing. A lot of users will never need a database that can do a million reads per second across a thousand machines reaching p50 1ms on 2 petabytes of data, but Cassandra can do that, and you don't have to build a whole sharding layer on top of mysql/postgres/redis or even install Scylla to get there.
Cassandra is commonly used as a sink for complex service telemetry e.g. mobile data from the carrier perspective.
e.g. you're willing to model your reads such that you dont need joins, and your sorting/scanning is based on the idea that you cluster data together, sorted in order you'll read it, into something called a "partition", which is how data is found within very large clusters.
RedPanda seems to be a better starting point to be honest - a lot of the market share for Cassandra doesn't actually need the extra tps of Scylla. Some customers do, but many many many do not.
Disclosure: I'm an engineer at ScyllaDB.
The ease of maintenance, similarly, is sold as easier due to reduced node count, which is perhaps an extension of performance but probably misunderstands (or ignores) that most people running large cassandra clusters have tooling that parallelizes most maintenance anyway, so the reduction in effort is sorta not that important in real life (if anything, having more machines gives you better blast radius behavior, consolidation onto fewer exposes you to larger percentages of loss/failure when there's inevitably a problem with the fewer, larger machines).
The real comparison, though, is missing in that link, because the real comparison is not performance. It's license. Nobody is running AGPL in prod unless they have zero IP worth protecting, so it's ultimately comparing OSS to proprietary.
(And similar disclosure: cassandra committer)
Thank you so much for this comment. I have come to a similar conclusion recently -- as long as the P100 latency is acceptable (e.g. all requests are served within 1 seconds), the only thing that matters is the TPS.
Back in 2016, I stumbled upon "How NOT to Measure Latency" by Gil Tene, and it really opened my eyes, especially the part about coordinated omission from benchmark tools.
Of course, after a few days, I also learnt that Gil was behind the Azul Zing JVM with the pauseless C4 garbage collector, and thus it makes sense for him to emphasize on measuring latency correctly.
At around the same time, GCP also boasted about having "consistent single-digit millisecond latency" for its BigTable offering, with Cassandra being the obvious target to attack. I was sold.
Then came Scylla, with a focus on maintaining low tail latency while running on fewer larger machines. I tested one of the early version with cassandra-stress, and the result was worse than cassandra. But I continued to follow Scylla blog posts with great interest.
Recently, I saw https://www.p99conf.io/, and something "clicked" when I read that it was sponsored by Scylla. Suddenly the hype of P99 seems to be wearing off. The video from Gil above is still correct, but I think it is applicable more to exceptional cases, e.g.
- the JVM enters full GC and do no meaningful work for seconds/minutes
- the InnoDB engine stalls for seconds/minutes due to a million different reasons
For normal cases, I haven't found the P95/P99/P99.9 metrics to be that useful. Instead, something like PostgreSQL/MySQL slow log threshold, where anything that exceeds a soft P100 target is logged, seems to be more useful.
Back to the article. If we only concern about TPS, under the "real-life" workload with Gaussian distribution, Scylla beats Cassandra by 2X. So that's it, 2X better performance on one hand, Apache vs. AGPL on the other.
(Of course that's an oversimplification. Scylla's shard-per-core architecture also allows it to avoid some silly single-threaded bottlenecks in Cassandra, but nobody wants to talk about those.)
https://www.scylladb.com/2021/04/22/on-coordinated-omission/
Also, yes, P99 CONF (https://p99conf.io) is sponsored by ScyllaDB, but we're very glad to have speakers from other NoSQL vendors like Couchbase and Redis, as well as folks from across the industry — streaming systems like Kafka (Confluent), Pulsar (Splunk), Redpanda (Vectorized). Plus storage systems like Ceph, Crimson (both Redhat) and Lightbits LightOS.
For the P99 CONF I recently conducted a quick poll of what people consider "acceptable" P99s. For some people it's <100 µseconds. For others, its <1 ms. And for some it goes all the way up to 1 second. But "acceptable" is use-case specific. An in-memory database or cache will have a very different expectation than someone writing data to SSD or even today, HDD.
33.3% of respondents expected <1 ms.
37.5% were okay with 1-<10 ms
20.8% were okay with 10-<100 ms
Only 8.3% were okay with 100ms - 1 second
(The <100 µsec was a "write-in" comment. But I am sure if I had included it, and we had a broader sample poll, you'd see it as a prevalent and vocal minority.)
I am totally a fan of Gil's work, but I think he over-dramatized the impact of tail latency, with the rather extreme example of "If a typical user session involves 5 pages loads, averaging 40 resources per page" from his talk.
Regarding coordinated omission, I am not totally convinced that an open-model system is what happens in production, and I am actually fine with a closed-model system when running benchmark, as long as p100 is acceptable. My goal is to achieve the maximum TPS without stalling, and I don't really care that much whether p99 is 10ms or 100ms.
On a deeper note, it probably says something about the user base that people who don’t run it in prod think the perf is really bad and yet there are very few people submitting perf improvement patches. Maybe most of the power users aren’t that worried about perf because they either know how to tune a JVM or they happen to size their clusters based on bytes on disk?
With scylla TPC approach, there will be one shard per core, and I think that shall allow more LCS compactions to run in parallel. Theoretically.
I have seen the workaround on a very large cassandra cluster. Meanwhile, I never thought to myself that "well, the 200ms p99 with g1gc sure is nice, but my users would be so happy if I can lower it to 20ms". So, if I am scylla, TPC would be my main selling point, coupled with real world examples, but somehow they want to focus on p99.
Also, Red Hat's Ceph's replacement, "Crimson"
Why is it faster? Is it algorithmic, or some neat trick, or just a much more efficient implementation?
And are they closely equivalent? Would you use one or the other for the same thing, or do they make different CAP promises?
Scylla is written in C++ (versus Java for Cassandra) and uses the high-performance Seastar[0] framework.
> And are they closely equivalent? Would you use one or the other for the same thing, or do they make different CAP promises?
Scylla claims to be a drop-in replacement for Cassandra.
seastar -> sea star -> C* -> Cassandra ;)
Scylla has a highly optimized low-level I/O path that largely bypasses high-level OS APIs and kernel services that Cassandra and other open source databases tend to use. This will typically generate an integer factor improvement in I/O performance if implemented well and makes a big difference for the kinds of write-heavy workloads Cassandra was built for. It requires taking strict control of low-level memory access and behavior, which (for better and worse) is the default case in C++. This code is intrinsically non-portable.
Additionally, there are some important classes of throughput optimization that are incompatible with garbage collection. In principle you can abuse Java to effect these optimizations but it is much easier to implement these optimizations in languages that don't have a garbage collector. If absolute performance is the objective, like Scylla, it is easier to do the implementation in a language that won't fight your intent every step of the way.
tl;dr: The performance isn't so much that it is written in C++ but that C++ makes critical optimizations relatively straightforward and economic to implement.
Can you share some more details about these optimizations? I.e. what they are and why GC tends to go against them?
Performance is adequate for Cassandra, so the community has (for several years) primarily focused elsewhere. It will be a priority again in future, but in the meantime with many huge scale users out there the community has focused on guaranteeing correctness and stability at scale. For example, the Harry[1] toolkit for validating huge databases, and an adversarial cluster simulator[2] for exposing distributed and other complex bugs. Also a huge amount of behind-the-scenes work that isn't so easy to call out.
The community is now focusing on expanding the utility of the database for these use cases. For example the recently proposed enhancement to bring state-of-the-art general purpose transactions[3] to Apache Cassandra.
[1] https://github.com/apache/cassandra-harry
[2] https://cwiki.apache.org/confluence/display/CASSANDRA/CEP-10...
[3] https://cwiki.apache.org/confluence/download/attachments/188...
[edit] disclaimer: I’m an Apache Cassandra contributor involved with some of the above work.
Definitely newer JVMs are improving things such as latency, and there are now some that are NUMA-aware, but in a JVM you are literally straight-jacketed from seeing the raw hardware you are running on. And that will impact to greater or lesser degrees your ability to take advantage of it.
When operating a database at huge scale surprising things happen, because everything that can happen will happen. So operators are interested in ensuring the database behaves well in these extreme circumstances. This isn’t specifically about horizontal scalability, though that is a necessary component.
To your point about vertical scalability, no doubt Scylla performs better here. However the details of your mentioned comparison are perhaps misleading, as Cassandra can happily exploit more than 16vCPUs before its performance materially plateaus.
While it’s true that the JVM imposes some restrictions, they do not translate to a difference in performance on the order of that claimed in this post. JVMs have also been NUMA-aware for some time. The main explanatory factor is relative investment and focus.
I suspect this talk is based on this post.