Revisiting 1M Writes per second
techblog.netflix.com
techblog.netflix.com
http://www.youtube.com/watch?v=c12cYAUTXXs
That's 'Billion' with a 'B': Scaling to the Next Level at WhatsApp
(that walk title was create before the acquisition and was mean to imply message count, after the acquisition it got a secondary meaning).
The one thing that is fascinating about it, is how small their team was compared to the volume and complexity of the operation.
I know I can divine it from the parameters to stress, but I have no idea if the row keys generated by different clients overlap and I don't know the default number of columns nor their size.
I think it's also important in this kind of benchmark to describe the distribution of access especially for a read intensive benchmark. Without that you really don't know what your are looking at. I am a fan of scrambled Zipfian.
That said, there's a previous benchmark linked to at the top of the post:
http://techblog.netflix.com/2011/11/benchmarking-cassandra-s...
The client is writing 10 columns per row key, row key randomly chosen from 27 million ids, each column has a key and 10 bytes of data. The total on disk size for each write including all overhead is about 400 bytes.
There are 3 replicas, so figure that in as well.
For what it's worth, 3.5m/yr is about 0.08% their revenue.
Microsoft do this all the time - its great publicity. It would be interesting to see if they (Microsoft) have something showing similar results...
Here's some performance notes on the latest (2.1rc4): http://www.datastax.com/dev/blog/cassandra-2-1-now-over-50-f...
Because there are many companies (including ours) who would strongly disagree with you about that.
285 nodes that are easily automated to create/destroy/monitor doesn't seem like a management pain to me, personally. Just depends if you have the write internal tools built.
I'm not saying I'd do it by myself. I'm saying with the write tools its doable and not unreasonable.
Though, I agree if you dataset is large enough and you need random access it's not going to help much.
The type of companies who could afford this are the types of companies who need 1M writes/second. Which are few and far between. And yes 10K is not that impressive but 1M is. And with Cassandra you could continue to improve that number just be rolling out more nodes.
Many companies need far in excess of a million writes per second. Basically, most machine-generated data sources, whether it is personal location data or any other kind of telemetry. Many companies that do not generate that data themselves buy and consume it. I know of companies doing over a billion writes per second.
Cassandra is pretty good for this type of thing among open source software but it is not nearly as efficient as it could be in terms of write throughput. If the storage engine is correctly designed, you should be able to drive 10GbE all the way through storage -- call it 1 GB/sec per node. However, that does mean you can't do things like mmap()-ing files; those interfaces are slow due to poor scheduling by the OS when the throughput is very high.
If you are doing it well, 3-5x throughput improvement seems to be average upside in my experience, which is huge. The scheduler behind mmap() simply does not have enough context about the workload to make good paging decisions leading to a lot of suboptimal or wasted I/O, and this is magnified when the storage I/O is under pressure. In principle, if you write your own I/O scheduler you can always make sure that the optimal I/O operation is executed at the optimal time.
How big was your total dataset. It's cheap to average a million writes per second if the dataset fits in RAM, or at least if the index set does, with the right database. It can be less cheap for a data set far larger than RAM, as for most databases write amplification becomes a significant problem.
While there is some write amplification it is less than most databases. It only takes few disks before the scheduler can get significantly more bandwidth out of the disks than a 10GbE network has to drive that activity, so there is extra capacity. The bottleneck on most server hardware is the silicon between the storage and memory if you are doing it right.