Looking with Cassandra into the future of Atlas
metabroadcast.com
metabroadcast.com
(In all seriousness, what did they expect? Have they looked at the MongoDB code? Do they seriously believe that the 10gen folks are smarter or better at solving problems than the masses of engineers Google and Amazon have thrown at this problem?)
That said, part of the reason Cassandra was attractive to me from the beginning is that unlike master-slave designs like MongoDB (or Bigtable/HBase, for that matter), a p2p design doesn't have the many corner cases around failover and recovery that complicate troubleshooting so badly. This is a primary reason Cassandra has had a very good reliability story since very early on.
What I really meant to say is that it's clear that the engineers behind Cassandra have done their research and chosen an extremely well tested design, while the engineers behind MongoDB seem to be completely winging it, ignorant of the literature and writing (based on my last examination) extremely amateurish code.
> Amazon, Google, Facebook throw mass engineers at problems Cassandra has overlapped with.
> There for, Cassandra is better than "database written by amateurs"
If you have hard evidence post the fucking evidence.
My evidence is "the technology that Cassandra uses"--that is, the SSTable/BigTable storage format and the Amazon Dynamo eventual consistency and distribution model--has been not only peer reviewed and published, but has been serving for years as the foundation of two of the busiest Internet websites in the world. That's real world experience orders of magnitude more extensive and robust than the anything MongoDB can offer.
Sure, it's got some cool big data stuff, but try doing any of those "maintenance" operations on other databases without ripping your hair out. For example, even bringing up a new MySQL slave is a huge pain in the ass, let alone doing something non-trivial like promoting a new master.
Thanks.
Even the developers who were resistant due to the different data model have completely embraced it.
One of MySQL's greatest strengths is the maturity of its ecosystem. Tons of existing automation, a relatively large number of skilled developers you can hire, multiple dedicated conferences, good free online docs, dozens of books, several top-notch consulting firms for handling the crazy edge cases.
I'll readily admit that MySQL has its flaws, but personally I wouldn't rank difficulty of maintenance operations high on the list... at least from a practical perspective, and relative to other comparable data stores used at a massive scale.
In my experience, no matter what databases you choose, if you "go big" you can't ever escape the need for robust automation and a decent amount of in-house talent.
But cluster elasticity and resiliency has been working really good so far.
Facebook never used modern Apache Cassandra, so when they tasked a team of experienced Haddop/HDFS engineers to build Messaging, it was natural for them to choose HBase. But the things they've had to do to deal with the problems in that architecture (e.g., sharding into "cells" to deal with namenode SPOF) [1] make me think that they would have done better with Cassandra.
[1] http://www.slideshare.net/brizzzdotcom/facebook-messages-hba...
With cassandra you choose where the extra reads (or writes) to ensure consistency occur as best suits your use case:
* At write time (write quorum) like a traditional replicated datastore
* At read time (read quorum)
* In the background (read one, with a non-zero chance of read repair)
This is quite nice to be able to bend Cassandra on a use case by use case basis (NB: these are not cluster wide approaches, for different columns / circumstances i can choose different patterns).
Talking about DynamoDB, it is certainly attractive as it scales, is managed by Amazon and all the yadda-yadda, but we didn't want to be lock-in to Amazon, and costs were another concern.
I agree with you though, locking into Amazon is kind of a drawback :-(
First, in order to achieve strong consistency in a distributed system with replication enabled, you will always have to write to multiple replicas: provided the number of replicas is the same, whether you're using a consistent-hashing based replication (Cassandra) or a distributed file system (HBase) it really doesn't matter much, as I/O latencies (both network and disk related) will always play a large role here.
That said, you can have different performances if you:
1) Change the replication factor: this is something you can do regardless of the actual system you're using, so it doesn't count.
2) Change the number of "acknowledged" writes/reads as a fraction of the replication factor: strongly consistent systems will always provide at least as much acknowledges as Cassandra's QUORUM, so no big win here.
3) Change the point where you acknowledge writes, either memory, disk or physical sync: this is about durability, and if you use different settings for different systems (i.e. Mongo VS Cassandra) it is apples and oranges.
4) Change your physical storage structure (provided same level of durability as explained above): different databases use different structures, Cassandra uses Memtables+SSTables, which I don't expect to be slower than HBase HDFS-based ones.
So in the end, Cassandra is at least as fast as other distributed databases, even when working in strongly consistent fashion.
That said MongoDB has a simple flag which allows you to choose how to handle writes. The app can 'wait' for writes to be applied to one node, all nodes, some number of nodes etc.
Dynamo (and thus Cassandra) is designed for write speed and durability. What is MongoDB designed for?
It has a per collection lock and is moving towards per row locking.
(ok, sorry for the obligatory Illiad joke there.)