Apache Kafka, Samza, and the Unix Philosophy of Distributed Data
confluent.io
confluent.io
For some contrast, take a language like python, which will allow you to create example programs to test out Kafka VS RabbitMQ VS zeromq. You will quickly find that one of the message bus daemons constantly gets in your way, fails to reliably deliver messages, and consumes a ton of system resources compared to the other two. Hint: It's kafka!
I really can't say enough bad things about kafka. Having used it and unfortunately been forced to implement it at a number of "big"(juniper, cisco, vmware) companies, it has been a horrible and disappointing experience for both the end-users and the developers, every single time.
tldr; don't use Kafka, use a real messaging bus.
I'm evaluating Kafka for a new project and it seems to be a perfect fit. I've contemplated building something from scratch in python, as my reliability and performance demands are pretty minimal. However, it seems that a lot of thought went into Kafka's design and it's feature set is perfect match for my problem. Specifically the unlimited buffering, log compaction and the ability to replay logs from arbitrary offsets.
If there are any viable alternatives to Kafka what are they? Bonus points if the JVM isn't involved.
However, for operations not involving filtering, i.e. when processing every line/record, you can increase the throughput avoiding IPC memory copies and parsing on every step (the copy + parsing could take an important portion of processing time).
> awk doesn’t know about the format of nginx logs
So the complexity of the implementation is a direct function of the dimension and heterogeneity of the data[-], but it still reduces to:
pgsql notify | client daemon | process | search engine
[-] high throughput is another matter
The question of how to build an "operating system for the internet" is really interesting. How can we make it as easy to write a networked program that can be composed with other networked programs, in the way that you can compose tail and grep?
Kafka and Samza aren’t it; they're cluster-centric, not internet-centric. They assume a single administrator, Kafka's publish-and-subscribe model is inherently unreliable in the presence of spammers or other sources of system overload, and spammers are unavoidable if you can't kick them off the system. (If you're going to run Hadoop in "secure mode", then it uses Kerberos for authentication; I need say no more.)
With regard to system overload, this is exactly backwards:
As long as the buffer
fits within Kafka’s available disk space,
the slow consumer can catch up later.
This makes the system
less sensitive to individual slow components,
and more robust overall.
You can reliably feed terabytes of data through a Unix pipeline consisting of a fast process feeding that data to a slow process (e.g. rot13 | some slow Perl script). Without backpressure, you have to buffer the terabytes of intermediate data, which will fill up your disk, causing your system to fail; or you can drop data, causing your system to fail. Systems without backpressure cannot provide this kind of composability reliably. (Similarly, systems that require you to name your data outputs in a global mutable-object namespace, like Kafka topics, pose obstacles to composability. This, more than anything else, is what limits the composability of queries in traditional SQL databases.)Don't get me wrong. I think pub-sub is great, and especially for loosely-coupled integration of different systems. I use pub-sub systems every day. I've even written a few. But don't make the mistake of thinking that they're the networked equivalent of Unix.
IPFS might be it. There's a whole panoply of other projects which, like IPFS, are trying to build a decentralized internet OS (MaidSafe, Ethereum, and so on) and probably sooner or later one of them will succeed.
One more minor quibble, which to my mind shows that the author wasn’t very interested in writing true things instead of false things:
you can pipe
the output of gunzip
to wc
without a second thought,
even though the authors of those two tools
probably never spoke to each other
David MacKenzie and Paul Rubin are in fact listed among the contributors to gunzip.I am amazed that these systems, Storm / Samza / Spark / even the new Flink are all Java / Scala based.
What other language and runtime platform provides:
1) Productive, high-level programming languages.
2) Portability to any OS.
3) High runtime performance for the server-side usage model.
4) Real shared-state SMP (necessary for memory and core efficiency).
5) Standalone, single-directory distribution of the applications and all library dependencies.
6) A broad, stable, mature base of well-documented, well-tested libraries.
The CLR has some advantages over the JVM, but it also shares most of the JVM's attributes.
As for using other languages - I played with Python on Storm but I very quickly found that I had to basically use Java because the entire community and documentation was java-centric (last I checked - a few months ago).
I may be being obtuse here - in fact I know I am - but my point is that there is a large community out there which is not Java friendly (logically or not) and so I lament the fact that all the advancement in this new, very exciting, field, is JVM-based. My dislike for Java was sealed by the Bloomberg terminal API, some of which's basic functionality is 5 to 6 dots deep. Seriously I had 80+ character function calls. This.that.that.this.this.finally()!
Perhaps we'll just have to kick in, leave our biases behind and start running with the JVM.
FWIW my use case is streaming fixed income (bond pricing) data analysis. Python/Numpy is running out of steam fast for me and we're using quite a bit of C. I basically come from the scientific computing set. We're applying advanced statistical analysis to pricing (not HFT - we're operating in the 5 minute thru 1-week horizon - not 5ns).
Your point about documentation is well-taken. We are trying to document streamparse + Storm usage from a Pythonista standpoint via our online documentation[3], e.g. here is our detailed Python API documentation[4].
[1]: https://github.com/Parsely/streamparse
[2]: https://www.youtube.com/watch?v=ja4Qj9-l6WQ
So which one is it? :-)
The underlying goal in Pythonic data-science is that one pushes as much of the computation into pre-packaged loops that have already been written in a close to the metal language (C, C++, Fortran). Well, that's as far as the intent goes, practice deviates from it by degrees.
This does not work quite as well in Java (or JVM) because JNI is just supremely god-awful. JVM semantics are overly strict, this over-specification kills optimization opportunities. Is heavy on memory, I would rather use the memory for loading more data than fill it with overhead. Finally there is this OOP culture that gets in the way.
That said, JVM is one of the most mature, and well engineered VMs out there (Java is another story), but not very well suited for number crunching, because there is more to number crunching than calling BLAS APIs.
That said, you can use these tools as complete applications without actually dipping your toes into Java (or even the JVM, other than running it). If you do want to write code directly against them, the choices (of which I'm sure you'll well aware) range from Clojure to Scala to Java to JRuby.
Scala is one ugly duck, but personally I find it to be a sufficient language for writing JVM-based code without losing my lunch.
But second, not quite a "platform" but I've been enjoying Frank McSherry's blog posts in which he builds up a complete distributed dataflow engine in Rust: http://www.frankmcsherry.org/
Somebody's downvoting me but I doubt they know anything about streaming and short-job support in gridengine.
We also use Kafka natively from Python via the other module we released publicly, pykafka[2]. It provides support for nearly all of the Kafka protocol, and we even have an optimized C extension module in the works that may even be faster than the JVM consumer.
My view on this is that infrastructure has to be implemented in some language. I don't lament the fact that Postgres was written in C or that Cassandra was written in Java -- I just use these technologies as infrastructure. Both have very full-featured Python bindings. The only thing missing on the Storm/Kafka/streams side are the full-featured Python bindings. But my team at Parse.ly has been hard at work building those in open source. Help us!
My understanding is that Twitter's rewrite of Storm (Heron) was mainly to make it work with Mesos as a resource manager. This is probably wise since Mesos handles a lot of the concerns I described above. But Mesos didn't really exist when Storm was written.
They chose to re-implement 100% of Storm's API in doing so, which perhaps shows that the high level concept of the framework has staying power even though you might benefit from a more advanced resource manager at 1,000-node scale. I wish they had reimplemented it as a competing implementation and actually released it as open source. But no, they decided to keep it 100% proprietary. So it goes. I guess once companies become a certain size, they turn their back on open source if it's not directly in their interest any longer.
We run it with 10-15 Storm nodes with 32 cores each and find it to be immensely helpful in this context, keeping each node in the cluster lit up to ~60% CPU utilization and plowing through 10K events/second on a Kafka topic, despite the fact that we are using Python and there isn't a line of code using threads or process management. People are often shocked we pull this off, since, in theory, CPython's GIL means you can't even run on more than one core at a time. But we write simple Python programs that run on Storm and utilize hundreds of cores at once, across multiple machines. And we get the Erlang-style "let it crash" / "fail fast" process supervision for free.
I have a more cynical view of the Heron paper overall -- if you're curious about that, reach out to me directly (@amontalenti on Twitter).
Pykafka looks nice! Any interest in getting it working with asyncio (or the 2.7 version, trollius)?
This is because, under the hood, it all uses ShellBolt/ShellSpout in the Storm layer, which is all process-based. But I actually find this to be the "purer" way to run Storm, anyway. Why deal with threads when you don't need to.
As Joe Armstrong, the creator of Erlang, once said (paraphrasing): "Processes are isolated environments for code where state can't be shared except through explicit messaging; threads are isolated environments where state is directly -- and dangerously -- shared. Why would you want threads when you could have processes?" Ofc, I realize, there are times when threads' lightweightness matters, but it doesn't in our case, and processes are certainly simpler!
As for your question, yes, we are working on async support for pykafka. The async producer is being worked on in this issue: https://github.com/Parsely/pykafka/issues/124 -- feel free to contribute, or even simply +1 as a vote of confidence!
Hi Andrew,
The Heron rewrite had more to do with what were at the time gross operational inefficiencies of Storm at very high scales, and problems diagnosing failures and bottlenecks. This is Storm 0.9 -- in the years since, I think Storm community in general and the Yahoo folks in particular have been working hard on addressing some of those issues, and a recent blog post from them indicated that some of the stuff we fixed in Heron is on their roadmap. Note that the Heron paper was published a year or so after the first Heron topology went into production inside Twitter. We had a very real problem that we needed to fix very quickly, and writing Heron was faster than making Storm work. Some of that was due to OSS challenges, some due to Storm's architecture fundamentals, some due to the specific people and backgrounds we had on the real-time compute team at that point.
The scales I am talking about are hundreds of nodes, not dozens, and an order of magnitude more messages per second. I am not surprised it works perfectly well for your use case (it worked fine while we were only putting tweets into it in 2013, as well -- that was about the size you are quoting, iirc).
Mesos not only existed when Storm was written, Storm ran inside Twitter on Mesos since before Storm was open-sourced. Mesos went into Apache incubator in 2011, while Storm did so in 2013. I'm not sure why you got the impression any of this was related to Mesos; it's true that we simplified a lot of operational complexity for Storm+Mesos by not writing our own Mesos scheduler, like Storm did, and just using Apache Aurora -- a decision that also meant we were able to use Aurora/Mesos clusters shared with other processes, which was nice; but that wasn't the prime motivation by a long shot.
The API was kept as a trade-off, to make migration of internal customers seamless. There are problems with the API, particularly around back pressure, but at the same time it was a straightforward API that got wide internal adoption, both in raw form and through Summingbird. So it made sense to evolve that part incrementally and deliver our internal customers the performance, observability, and reliability wins first, without having them rewrite a line of code, and tweak APIs over time to address the above-mentioned issues.
I'm biased of course, but I think our track record with open source contributions is still pretty good -- Scalding, Parquet, Aurora, Mesos, Finagle, Zipkin, and many more major projects with very wide adoption in the industry, which we still very actively contribute to and continue to evolve. At the same time, there's definitely a cost to open-sourcing projects, and those decisions get made on a case by case basis.