Scaling Riak at Kiip
basho.com
basho.com
You're no doubt wondering what the issues were; they come 20 minutes in, and were (paraphrased):
* At scale in production, adding a new node took days to complete all the handoffs; they recommend adding new nodes as soon as it's looking like you need them, rather than waiting until you're redlining.
* 2i is slow, especially in EC2; a straight KV "get" is milliseconds-denominated; 2i index queries were taking multiple seconds. Use 2i, they say, but in background processes.
* Javascript MapReduce is slow; this is well known. They confirm Erlang MR was adequate.
* As the LevelDB keyspace grows, there's a stepping function in latency; 5ms, then 15ms, then 25ms; the solution is to add nodes. (LevelDB is Google's KV store, a new option for Riak, required if you're using secondary indexes).
* Riak Control didn't work for them over low-latency connections.
* Once, a Chef misconfiguration left the whole cluster flapping on and off, which corrupted the cluster; they recovered with Basho support. Be careful about adding and removing nodes rapidly.
* Similarly, flapping a single node caused the cluster to get into a state where it wouldn't converge again; the cluster worked but no nodes could be added until they (presumably?) restarted it.
Thanks anyway.
Last night I just finished replacing 6 maxed out medium instances with one $100 box from SoftLayer. :/
edit:
Softlayer: Intel 4x2.40GHz, 2 GB RAM ECC, $160/mo Honelive: Intel i7-2600 16, 16 GB RAM, $120/mo Kimsufi.ie: Intel i7 4x 2(HT)x 2.66+ GHz, 24 GB, $60/mo
That you'd encounter all these things at a 25mm daily ops level is pretty odd, though.
Yet, they went and scaled horizontally with Riak and experienced pain.
Their opinion that it did not make sense to have to horizontally scale the "relatively small" number of ops they were sending to MongoDB is certainly their own, but then they horizontally scaled with Riak anyway and boasted about their 25MM ops per day scaling ... which, averaged out, is only about 280 ops per second.
In short, it was far from an apples to apples comparison.
That's a bit of a headscratcher. What is happening during those 'days' and what is the primary limiting factor?
I keep meaning to get into Riak, but then stuff like this where the system has crazy moments that are impossible to coherently reason about keep popping up.
The operational challenge I infer from this is that they had waited to add that node until they really needed it, because their expectation was that adding the node would get them quick relief to their scaling issue. Instead, they got relief a few days later when the node was fully integrated.
Solution: don't wait to add nodes until the last minute.
We have a fairly small cassandra cluster in production, serving over 50x the volume they mentioned in their talk, with good latency (real-time bidding) and not-too-painful operational footprint.
Just like with anything in technology ... there is going to be pain as you escalate the level of complexity and what you are trying to accomplish.
The issue was that underlying architectural decisions in MongoDB ended up biting us and limiting us rather significantly. It could be argued that this is because MongoDB is a rather new piece of tech (I disagree, I think MongoDB is fundamentally flawed, but it doesn't matter in this argument).
Because of our experience, going with the standby IS the best choice, until you _need_ something else. Riak was a change necessitated by its fast growth such that horizontally scaling was necessary when you're in an environment such as EC2. Ignoring MongoDB specifically, the IDEA of MongoDB simply isn't correct here, a Dynamo-style K/V store is the correct option, and Riak happens to be a fantastic one.
I love the database but it is not optimal if you have a fluid schema like many use cases have today.
With all due respect to the Kiip engineering team, this wasn't a strong case for using Riak over MongoDB ... but rather the general pain that a engineering team feels when horizontally scaling in the cloud.