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.