Kafka, on the other hand, when you write a message to the broker the broker writes it immediately to disk queue rather than holding it in memory. But isn't that slower? No, it's not, because it's in page cache, which is managed more efficiently than garbage collected memory. Then, when consuming, rather than keeping metrics for each individual message being received, consumers simply have a log position -- they periodically commit, which tells the broker that all of the messages until that point have been consumed. If they never commit, eventually another consumer will get those messages.
So basically, it scales a ton better because you're just doing scads of sequential I/O with occasional commits, rather than tracking a bunch of messages in memory individually (which in theory should be fast but causes GC problems).
EDIT: should add that morkbot had a great link too:
https://news.ycombinator.com/item?id=6874607 http://www.quora.com/RabbitMQ/RabbitMQ-vs-Kafka-which-one-fo...
:-)
Pre-0.8, if a machine fails you lose all the data on that machine, only the lower durability levels were available. It guarantees at least once processing, while other queues generally make stronger claims. etc.
It's still very good at what it does.
What you're talking about is failover and fault tolerance, which are greatly improved in 0.8 with the addition of replication.