The trouble with Cassandra as an object storage metadata database
blog.min.io
blog.min.io
Okay, I'll just not bother with partition keys, right? Except partition keys determine which "partition" your data goes in. So if your only partition key column is "day_of_week", then you have 7 partitions. New problem: Your partitions need to be < 300 MB or cassandra falls over dead.
In fact you'll soon realize that over the lifetime of your table, keeping those partitions under control may end up forcing you to use various date tricks like putting year/month/day in as "artificial" partition keys.
Of course if you put everything in the partition keys instead of clustering keys, let me note again that you have to put each partition key column in every where clause, in which case you may have trouble querying for batches of data.
Furthermore, when you do put clustering columns in your where clause, they have to be in order declared; so if your CC's are a, b, and c, then your where clause can use (a), (a,b), or (a,b,c); but if it has b, it has to have a, and if it has c, it has to have a&b. This is because storage is hierarchical. (Same rules for "order by", btw) (and no there's no "group by")
This is when you start realizing: Oh, you mean it's really not like "SQL without joins". No, not even close.
HashMap<PartitionKey, OrderedSet<ClusterKey>>
So you can lookup data by the ParititionKey, then perform range queries on the ClusterKey to filter data belonging to that partition.
If your access patterns look like that, Cassandra could be a great fit.
For anything else, you probably want to consider a different data store.
Though HashMap should be just "Map" and the value of the top level map is not a set, it's another map:
Map<PartitionKey, SortedMap<ClusterKey, Record>>
C* is essentially scaled/automated MySQL Sharding and blob storage.. except you pay a huge cost for any server side filtering.
The Partition Key would be the equivalent of your MySQL shard id. ClusterKey is your primary key. if you treat everything like PUT $paritionKey, $primaryKey, $data and GET $partitionKey, $primaryKey you get a slow and massively scaling redis. If you do anything else with it you'll likely start regretting your choices and looking for a replacement database.
So for your "primary" data store for a type data, you might use the date for the partition key. Then if you do a lot of day-of-week queries, you'd store a second copy grouped by year-month-weekday, and query as many months as you need, and combine them. If you are collecting data for several sites, have a copy that's stored by year-month-weekday-site. And set up your software so it's easy to say "write in these 5 places, keyed off this combination of values." And because it's in all these denormalized places, you can forget updating anything atomically. So you go to append-only mode and if you have an "update" you record that as a separate entry after the first and your application's understanding of the data store includes coalescing the events and applying the update. (Maybe if you need it you have a separate process coalescing updates after the fact, but in the background, non-critical-like.)
And you're right. This is very much not SQL. It's also super obnoxious and support for it in extant frameworks is pretty minimal (they tend to be rather myopically SQL-oriented, ill suited to take advantage of this model). It also uses lots of storage, on the premise that disk is cheap. But what all this accomplishes is crazy horizontal scalability, and when you make any one query, then all the data is right there on disk lined up neatly, so it all gets streamed out in a predictable amount of time. For a few select applications, this consistent scalability is worth the pain.
people work around that by storing the data in something like S3 and then keeping the handle in cassandra. Yet another hack on top of another hack.
You'd do the same "hack" there.
https://www.scylladb.com/2020/12/15/scylladb-blog-scylladb-d...
> Your partitions need to be < 300 MB or cassandra falls over dead
Why is that? Is that some kind of engineered-in limitation?
It's getting better but likely will never reach the performance potential of Scylla which is Cassandra reimplemented in C++ with a much better sharded architecture.
Is that a typo? 300MB sounds ridiculously low..
Other than what the article described, I can also add:
1. It has a steep learning curve, but you do get to see the advantages while you learn it. But then, everything comes crumbling down.
2. The setup is a pain locally. Then it is a pain to set it up in prod and manage it. The tooling itself feels very unfinished and basic.
3. No querying outside primary index on AWS Keyspace if you want it managed. Also, any managed variants are EXPENSIVE. I mean, every database is fast if you only query by the primary index so why pay extra?
It is just not worth it. For example, we winded up using MongoDb and it turned out to be fast, scalable, had mature tooling and we can keep tons of event related metadata in it and it is easy to manage and doesn't cost a fortune.
I haven't really worked in this space for a couple of years so I don't know if the cloud offerings have already completely matched Cassandra's features and robustness.
Of course, I have been totally unable to determine how they merge rows/partitions without cell timestamps. It's a black box.
I was just on a "Keyspaces" meeting where the sales dude basically described dynamo billing, dynamo provisioning, feature shortfalls obviously due to dynamo, but would not admit it was dynamo.
It was bizarre.
We ended up going with ScyllaDB, which is a drop in replacement for Cassandra. It’s written in C. Much easier in resource demands and we didn’t have to deal with Zookeeper directly.
The parent comment almost seems generated by AI.
The Minio team are manifestly excellent engineers, but insightless posts that contain subtle misunderstandings of CAP do nothing to showcase that competence.
I'm leaving my original reply below because what the hell.
~~~~
Yes, exactly.
Think of it this way: you don't actually have any control over partitioning, therefore partitions are a given. So CAP is expressed as such: given a partition, choose between consistency and availability.
EDIT: I personally find PACELC [0] less confusing, and more nuanced than CAP, but they basically say the same thing.
https://github.com/minio/minio/blob/master/pkg/bucket/object...
https://github.com/minio/minio/blob/master/pkg/dsync/drwmute...
And this is their test:
https://github.com/minio/minio/blob/master/pkg/dsync/dsync_t...
I just browsed quickly but it is littered with Amazon Simple Storage Service hardcoded bits like this:
https://github.com/minio/minio/blob/master/pkg/bucket/object...
There is not a single document that I can find that discusses the MinIO architecture. I guess "MinIO is a simple wrapper around S3 with a homegrown distributed state tech using NTP and it is 'fast!'" does not make for a sexy doc.
The column "trouble with OSS distributed DBs" without Jepsen tests has probably already been written. There should also be one about "competitor's mature product bashing blogs are HN clickbait".
My memory is letting me down here on specifics - in my defence it's been a long time since i answered questions on stack overflow about Cassandra. I absolutely loved working with the product even if the sharp edges cut me a few times (ultimately my own fault not Cassandra's).
- they have Cassandra committers (most end up going to Apple iCloud)
- they have mgmt. support to make it work
- they are/were on tokens for years, not vtokens, since their internal tools were token-based
(They started with Oracle Enterprise (and a little MySQL), tried Mongo briefly, then decided to make Cassandra work for them.)
Source: original Netflix Cassandra team member.
Your metadata database needs to be the fastest and most reliable store out of everything. It can't be eventually consistent without partitioning your datastore. Even then you'll end up partitioning your data neatly into the same failure zone.
Cassandra has basically one usecase: high volume writes, with a few batch reads.
Cassandra is not really optimised for high reads.
Most of the time postgres will do fine.
For something that scale horizontally, I would probably recommend to use something like FoundationDB for this use case [^1].
A transactional Key-Value store is exactly what you want for this kind of use case.
While you can vertically-scale a single PostgreSQL (with failover/ha) server to 50+TB of NVME and grow until you can throw other engineers at the problem.
Example: wasabi.com uses mysql for metadata.
Yes. It's also a good choice for horizontal scaling, but only if the other DBs won't do.
That's how we use it in BigCompany(tm).
Cassandra holds (almost) all of our data. Most of it comes from batch insertions, but more complex (and lower volume) data comes from various REST APIs we expose and other teams use.
Most reads (60%) done on it are boring batch reads: analytics, reports, deltas, etc.
For more complex reads (30%) we copy relevant rows to Hive tables, then perform queries there without dealing with partition key shenanigans. So, batch with extra steps.
Then anything that requires random reads (10%) is read from ElasticSearch:
1. Coming from REST backend if it's time sensitive.
2. Coming from a batch that reads from Cassandra if not.
The article listed many known Cassandra characteristics and cited them as limitations. However, it all depends on use cases. There are no file system that works for all cases, and not all of them needs ACID, CA vs CP, etc. The rest points are not convincing either. They are related to how to design the data structure better.
Actually, SeaweedFS can use many other database/KV stores as the metadata DB. The list includes Redis, Cassandra, HBase, MySql, Postgres, Etcd, ElasticSearch, etc. https://github.com/chrislusf/seaweedfs/wiki/Filer-Stores
I did find one drawback for Cassandra as the metadata store though. One use case is that the customer uploaded a lot of zip files to one folder /tmp, unzip them, and then moved to a final folder. The rate is about 3000 files per second created and then deleted. Being a LSM structure, the tombstones quickly pile up and the directory listing was slow.
The solution was to use Redis for that /tmp folder, and still use Cassandra for the rest of folders. With Redis B-tree structure, the creation and deletion are cheap.
So it is all depends on use cases.
Pretty sure redis has no b-tree but hashmap. But redis has to do no compaction at all since it can directly delete the stuff on disk.
> The rate is about 3000 files per second created and then deleted. Being a LSM structure, the tombstones quickly pile up and the directory listing was slow.
The problem here also is the gc-perdiod. Since the database has async-replication, it needs to keep deleted rows for some days so downed replicas can replicate. So they stay in disk/memory for some days actually and don't get compacted. While in a CP system they would be.
You may think "it will be cached in ram, because it's small", yes, but then you'll end up querying many nodes just for metadata queries.
Yes it's nicer to manage only 1 system, but in big scenarios it's probably better to separate.
You can have 50+TB of nvme in 1 server, so your metadata layer probably doesn't need to horizontally scale.
Imagine if you lose some objects (because you lost some replicas). You won't even know WHICH objects you lost, because the metadata is gone together with the data.
Having separate, you can add a 5-replicas to metadata to be even safer compared to the usual 3 replicas.
[1] https://www.duetpartners.com/why-is-premature-scaling-still-...
Has anyone ever used min.io for production stuff? What are the pros & cons over vanilla s3?
Pros:
* Less platform-dependent. By self-managing it, we can also run deploys to GCP / Azure / Scaleway / other providers without writing a separate adapter for e.g. Azure Blob Storage.
* Python API [1] much more pleasant to use than boto3 (and can speak to normal S3). It doesn't do everything that boto3 does, but it supports everything we need (e.g. pre-signed URLs).
* minio server itself supports a large chunk of S3's functionality (e.g. SELECT API / AssumeRole / bucket versioning)
* Don't pay per request and for egress: this was a big deal since people might want to download large amounts of data from us (or make a bunch of small requests to download/upload a subset of data).
Cons:
* Have to manage own infrastructure. We run it on managed VMs so it's semi-managed, but we still have to provision block storage, set up backup policies etc.
* In a similar vein, scaling and availability all have to be DIY [2]. We haven't run into situations yet where Minio would be the bottleneck, but it might be something to keep in mind.
* Obviously not as seamless: you don't get things like Glacier or integration with other IAM.
[0] https://www.splitgraph.com/
[1] https://github.com/minio/minio-py
[2] https://docs.min.io/docs/distributed-minio-quickstart-guide....
Mostly that you aren't bound to AWS IMHO. You can run it on-prem, or in a cloud provider with no object storage service.
The main pro and con is that it is self-hosted.
Somewhere else in the comments: "Yeah, we went with MongoDB."
Chances are they absolutely don't have the headcount to properly manage it and were just filling out the portfolio they wanted to offer.
The costs of operating the cluster are ~5K/month (that's what our service provider charges us for 24/7 ops). I consider this a scam since averaged maintenance costs are perhaps ~1 hour per month.
(normally peaks at 40.000 writes/sec, 500 reads/sec, 200GB)
An AP system may be a reasonable choice for an ODB if your added layer of transactionality is exploiting a domain specific characteristics that allows for fast-path, efficient, distributed transactions on top of a bare object level AP system (such as Cassandra). Since this is all very niche and technically demanding, it is almost always a better choice for the general team to choose a mature CP database as backing for metadata-rich objects and/or domains that demand transactional support for systems of record.
the team am in uses it, but after many times asking why it was chosen since it seems a poor fit for our uses cases compared to a relational DB, the only justification I was given is a that the cluster is easier to maintain for our ops guys.
Another reason is for "resume driven development".