Cassandra: Daughter of Dynamo and BigTable
insightdataengineering.com
insightdataengineering.com
[1] - https://www.facebook.com/notes/facebook-engineering/the-unde...
edit: Added link to facebook post on hbase vs cassandra
It really is the easiest database by far to scale that I've worked with.
At least Google's counterparts seem to be written in C++. Such as BigTable and GFS. Presumably also Spanner is C++.
In addition to C++, especially considering its less good safety/security record, Rust, Golang and Nim should be interesting alternative, safer, implementation language choices. In contrast to idiomatic Java, those languages provide significantly higher CPU cache hit rate for internal data structures due to no value boxing [1] and ability to reliably reduce problems like false sharing [2].
1: http://www4.di.uminho.pt/~jls/pdp2013pub.pdf (These issues can be worked around in Java by abandoning object orientation and instead having one object with multiple arrays (SoA, structure of arrays). In other words, not List or Array etc. of Point-objects, but class Points { int[] x; int[] y; ... } that contains all points.)
2: http://mechanical-sympathy.blogspot.com/2011/08/false-sharin... (False sharing performance issues)
And I don't understand why anybody should care about CPU cache hit rate. The bottleneck is always going to be in the I/O pipeline. And Java is faster than C++ and vice versa in various situations.
> Rust I understand provides some nice semantics for thread safety but these exist in Java as world.
How do you get Java to fail compiling if thread safety constraints are not met? I'd be interested to try it out!
One should definitely care about cache hit rate, because it significantly affects runtime performance. There are just 512 L1D cache lines per CPU core.
I/O is the bottleneck? It is becoming less so, one of the few areas where there's actually some nice progress happening. PCIe SSDs are up to 1.5 - 2 GB/s (=up to 20 Gbps). More and more servers have 10 Gbps or 40 Gbps networking.
Sure, Java is faster when it can use JIT to prune excessive if-jungle, aggressively simplify and inline and adapt to running CPU. But memory layout control is where Java is rather weak. The problem is getting only worse, because the gap between CPU and memory performance is only widening year by year. Memory bandwidth is increasing slowly and latency hasn't improved for a decade.
C++ is going to be always faster especially if specialized to certain machine and use case. C++ is also going to win by a large margin when there's auto-vectorizable code or heavy use of SIMD-intrinsics. 10x is not unusual, if the problem maps well to AVX2 instruction set.
http://damienkatz.net/2013/05/dynamo_sure_works_hard.html
TLDR quote:
The Dynamo system is a design that treats the probability of a network switch failure as having the same probability of machine failure, and pays the cost with every single read. This is madness. Expensive madness.
> Network Partitions are Rare, Server Failures are Not
Network partitions happen all the time. Sure, the whole "a switch failed and that piece of the network isn't there anymore" doesn't happen a lot, but what does happen a lot is a slow or delayed connection, or a machine going offline for a few seconds.
This is especially true for cross-datacenter rings across the public internet.
For WAN replication we have a cross-datacenter replication which works on an AP model.
I was simply trying to point out that while you may have a very good argument as to why Couch is better, the network partition argument is not sound, and you may want to look for a better argument to make.
I'm personally against single masters because they are SPOFs. With a master, at some point there needs to be a single arbiter of truth, and if that is unavailable, then the system is unavailable.
It's that these are fundamentally very different databases with different trade offs. You can't just take one and adjust some API calls and expect things to work in a similar way. It only confuses people when it's quietly ignored and others assume that since it wasn't pointed out to be wrong that it must be the same thing.
I've had far too many conversations with people who use Couchbase that can't tell the difference that I would say that it's just general confusion. It's lax work on Couchbase's part and a thorn in the Apache CouchDB project that there is no effort to help clarify the fact that they are indeed independent and now very different databases.
This is especially relevant when you need to do these things because of unexpected load increases or the loss of hosts in your cluster.
If there is a network partition, however, there is no need to move the data since it will eventually recover; moving it would likely make the situation worse. Cassandra never does this, and operators should never ask it to.
If you have severe enough network partitions to isolate all of your nodes from all of its peers, there is no database that can work safely, regardless of data shuffling or consistency model.
You can kind of get this with Cassandra if you write with CL_ANY and read with CL_ONE but hinted handoffs don't work so great (better in 3.0?) and reading with ONE may get you stale data for a long while. It would be nice, and I don't think there's any theory that says it can't be done, if you could keep writing new data into the database at the right replication factor in the presence of a partition and you could also read that data. Obviously accessing data that is solely present on the other side of the partition isn't possible so not much can be done about that...
You could then lose an entire datacenter (1/3 of the machines) and the cluster will just keep on running with no issues.
You could lose two datacenters (2/3s of the machines) and still serve reads as long as you're using READ ONE (which is what you should be doing most of the time).
You're susceptible to total loss of data since at any given time there will be data that hasn't been replicated to another DC and you're OK with having inconsistent reads.
That works for some applications where the data isn't mission critical and (immediate) consistency doesn't matter but doesn't for many others. I'm not sure what exactly NetFlix puts in Cassandra but if e.g. it's used to record what people are watching then losing a few records or looking at a not fully consistent view of the data isn't a big deal...
Data is only typically shuffled around when a replacement node is introduced to the cluster.
If you have multiple independent network partitions that isolate all of your RF nodes, then there is no database that could function safely in this scenario, and this has nothing to do with data shuffling.
Even VMs on more statically allocated clouds like DigitalOcean and AWS will experience small, constant blips that affect your whole stack.
What annoys me in particular is that these blips affect everything. Every app needs to fail gracefully, be it a PostgreSQL client connections, a Memcached lookup or an S3 API call. The fact that such catch-and-retry boilerplate logic needs to built into the application layer, and every layer within it, is still something I find rather insane. It leaks into the application logic in often rather insidious ways, or in ways that pollutes your code with defenses. Everything has to be idempotent, which is easy enough for transactional database stuff, less easy for things like asynchronous queues that fire off emails. Erlang has already provided a solution to the problem, but I suspect we need OS-level support to avoid reinventing the wheel in every language and platform. /rant
It's a heavy price, sure, but in return you'll be able to round off the tail in most scenarios. This is a trade off that many don't properly consider when they try to focus on that mean performance while they miss part of the point of the dynamo model which helps provide better guarantees about how things perform in more cases, including highly tuned clusters with very few major failures.
Note that Cassandra doesn't actually handle eventual consistency properly, and has weird corner cases a result (e.g. it's infamous "doomstones"). As an immutable data store it works very well, particularly when you have a high write load.
Not sure I could work in a one product company where you have to eat your own dog food in all situations even where it really is not suited for.
IMO, the only way use, or former use, is interesting is in understanding why a specific company moved to, or away from, a given technology.
Specifically, what was the original use case that was the basis for original use?
Why did the company choose to change technology? Did the use case change? Did the technology fail to satisfy the original requirements? Did a new technology with substantially better capabilities emerge? And so on?
Simply stating that company X uses Y (or company X no longer uses Y) does not provide a lot of information that other companies can use as the basis for their decision.
Of course if your data change often and very quickly you shouldn't be using Cassandra. You'll end up tombstone death. Deletes are not real delete, they're soft and just have a timestamp that eventually will be deleted (tombstone). Every delete creates a tombstone, if you're hashkey/column key have 50 tombstones, it have to go through those tombstones before getting the values. The reason is some trade off for faster write but shitty read if you update a lot of the same key.
Overall, one way of thinking of it is a souped up hashkey db that can do some relational queries which is much more than the usual hash key noSQL. Seeing how it's a hash key type database you can see the trade off of Cassandra and it's siblings versus say MongoDB or CouchDB.
I used it for time related stuff that has simple data. It's mostly immuatable so cassandra was perfect for it. Example is a daily tv show release date.
Time series databases like KairosDB are quite good, but for simpler data structures and something describable as a metric. Also you may face issues introducing relatively less known software, and get locked in, in your company.
I think Riak is richer in this area than Cassandra, because if the inconstancy can not be resolved, Riak can keep both versions and let you deal with it at the application level.
See also: https://aphyr.com/posts/294-call-me-maybe-cassandra