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).