Sharding a database can make it faster
stackoverflow.blog
stackoverflow.blog
Distributed relational databases have made tremendous progress in recent years and you have many compelling options now if you need it. For MySQL there's SingleStore (memsql), Vitesse/Planetscale, and TiDB.
Let's say you add a comment to a piece of content and then refresh the page. With a naively-implemented CRUD app, your comment disappears if your request's read query was routed to a lagging replica.
This gets even trickier when multiple datacenters in different geographic regions are involved.
Just about every other major database also supports replicas.
That's replication, not partitioning, and involves scaling up every instance to be big enough to hold the entire dataset. You must scale up first before you can scale out.
While modulo (along with random allocation) has an instructional value as an introduction to partitioning, it's best followed by a warning when it comes rebalancing: it violates a basic requirement for most partitioning schemes (with the exception of distributed partitioned logs), that is, rebalancing better move as few items as possible.
At any rate, the way I understood it, the article described partitioning of the data by `key mod k`, k being the current number of nodes.
Maybe I'm just getting salty, but why is there filler in this article. This isn't a recipe for garlic bread.
My mom was really my hero, because even as she took her fight for data sanitation out into the world she never lost site of building pure and clean relational engines from the ground up.
Really, my mom wasn't just my hero but a hero to all the kids in the district. Heck, almost every year we would spend days on the playground instead of the hated classroom thanks to her ingenuity!
But the best part of any day was when I was down in the sandlot with Butchie, Achmed, and Patty and my mom's voice would call down from our open kitchen window where she was baking and optimizing throughput "Bobby Tables, please accept these cookies I made for you and your friends in order to have energy for playing!"
That's why I am happy to share with you my mom's recipes for cookies and secure storage of cookie related acceptance data.
on edit: should be lost sight not lost site but will let it stand because hah.
Things are getting complicated at work. The app he's been working on for the past few months turned out to be a hit, new subscribers are appearing at every moment; in short, business is booming! They even got a column in the local newspaper, getting high praise for both form and functionality.
A problem presents itself when too many people are trying to use the app at the same time - it looks like the database software can not keep up with the demands of this immense crowd. An impromptu meeting is held, and Colin is tasked with coming up with a solution. Getting closer to the front door, he thinks to himself, "How can we handle this incredible volume of traffic when it reaches the database cluster?", but his thoughts soon drift off again.
If you truly need to shard, it's best to move deliberately and potentially roll your own at the application level. Your understanding of the infrastructure is far more important than sharding for the sake of sharding.
I guess it has to, but whats strategy there, if its mostly about "slicing and dicing" in every way possible.
Maybe recent vs. historical data, or some known customer segments, extracted from previous analysis?
Some warehouses like BigQuery force this as a numeric key (dates by default) while others leave it up to you. Partitioning can also use multiple columns to help with more even distribution and/or better query performance.
But you can't implement either if you haven't logically partitioned your data first.
Unfortunately, cross-table joins don't shard or scale well unless the data is well suited to a large scale sharding (like "customer" or "region". You'll probably start doing "data lake" type frameworks for join calculations across shards.
If you are sharding a database due to scale that modern big machines can't handle, you'll probably want to move your largest tables to cassandra or dynamo or other purely scalable approaches.
"fastest-developing verticals" ... sigh.
For many problems Cassandra doesn’t scale. For many problems log structured merge trees don’t scale.
- Super DELETE-heavy workloads. Tombstone cleanup can hurt.
- Complex multidimensional indexes. Things that you'd use multiple semi-overlapping compound indexes on a single table for in a traditional RDBMS.
- Operational pain of operating large Cassandra clusters. Recovering from multi-node failures while under high write load can be nail-biting as you watch the replication progress and the hint expiration "death clock" race for completion.
- Even semi-frequent scans over huge (GBs in size) tables can cause unexpected hiccups in other ongoing CRUD operations. If the scans hit enough nodes/partitions, the impact on other workloads can be pervasive.
- Emulating transactions (with an optimistic lock, full-cluster reads, and a cleanup program that expires interrupted transaction records) at scale can get costly. If you need this, cassandra was probably the wrong choice in the first place, but you asked for limitations and I have seen this one a couple of times.
Frequent DELETE: True, unless your data is time-windowed. The extension to this is frequent updates to the same "cell" of data, which also overloads the compaction and read paths. DELETEs have a big toll on relational databases at scale as well though, although SSDs mitigated a large part of those. DELETEs are very hard on B+tree indices without LSM features, and I think I've read that compaction issues exist in Postgres under heavy deletes: all those empty cells need to be cleaned up eventually...
Complex Indexes: Well, you can just roll your own view table, although I don't know your exact use case. The original (heavily downvotes) comment about B+ trees is that SQL users are used to throwing a B+ tree index onto tables. You can't do that in cassandra, and I'm not aware of how you would do it in other distributed scenarios that aren't full-dataset-replication strategies. From what I've read, a B+ tree scaling and many other index scalings aren't Cassandra problems, they are fundamental CAP problems.
Operational Pain: You are mostly right (Cassandra is a cranky beast), but once you get to heavy distributed data, it's all a PITA. At least Cassandra has some node loss tolerance. In a lot of times I hear people complain about Cassandra I ask, "well, did the database go down"? And they generally admit it stayed up.
Scans suck: Yup. If you need scans, do a separate datacenter for the scans. On a sufficiently large distributed dataset, all scans / "big data queries" are approximations of current data state: there's no meaningful/practical way to tie to a transaction window. Is this a cassandra problem or a fundamental CAP problem?
Transactions: Compare and Swap is about as good as it gets, and it's hard, and there's no way to span tables. Does anyone have a solution mathematically for sharded distributed data that isn't? As in is this a Cassandra problem or a fundamental CAP problem?
I agree this article is low quality and often incorrect, but I don't see what the choice of data structure has to do with it. Sharding is orthogonal to data structure.
For example, Facebook migrated their largest db fleet from sharded InnoDB (B+ trees) to sharded MyRocks (LSM). Still sharded either way.