Flying faster with Twitter Heron
blog.twitter.com
blog.twitter.com
Most notably, I couldn't determine whether the system is stateful, or what kind of guarantees (if any) are provided regarding stateful processing.
For example, the paper says Heron is used to compute real-time active user counts. That implies that the system needs some way to keep track of and "remember" how many unique users it has seen in the last N hours or whatever. How does Heron model this state and how does it guarantee (if it does) that a crashing node will not lose its accumulated state?
In my experience this is the hardest part, by far, of stream processing, so when I see any work in this area it's the first thing I am curious to learn about. A system that guarantees strong consistency (ie. accurate counts) even in the presence of node crashes is way, way harder to get right (and a lot more expensive, resource-wise) than one that assumes it's ok to lose a little bit of data.
It looks like Heron implements only at-most-once and at-least-once semantics, so maybe that is my answer there. You need exactly-once semantics to get robust and reliable answers, and you need to guarantee that state changes are atomic with the exactly-once semantics.
Of course some systems are ok with their output degrading a little when nodes crash. It's not the end of the world if the active user count is a little off. But beware of tolerating this too much -- the bad thing about allowing data loss is that it tends to come in storms (no pun intended). Once something is going wrong, the answers can be way off. The error is not bounded in most cases I've encountered.
The industry term for this approach is lambda architecture (http://lambda-architecture.net/)
I am looking forward to taking Heron for a spin.
Storm is a much better name than Heron, IMO.
This could be very good thing for Apache Storm depending on how Twitter handles it.
Just to clarify, the Storm version mentioned in the paper and blog post is not an official Apache release and doesn't include many performance improvements included in the newer releases of Apache Storm. There are a lot, and many more on the horizon.
That being said, the performance numbers look impressive, even though there is no way confirm those since no code or benchmarks have been published. IMHO, until that happens, there's not much to see here (not that I doubt it -- I'd just like to see proof/code).
My hope is that Twitter is dedicated to the projects it has open-sourced, and this is not a case of NIH, but rather an honest effort on Twitter's part to contribute back to the open source community.
@haberman:
Storm implements exactly-once processing through a higher-level API called Trident, that I like to call Storm's "Streams API" since it's not unlike Java 8's Streams API (and largely inspired by Cascading). Trident processes data in configurable micro-batches, as opposed to one-at-a-time, which gives it an advantage in terms of throughput, but at the cost of latency. Trident topologies "compile" down to Core Spout/Bolt topologies (The Trident API has a planner implementation that figures that out -- not unlike an SQL query planner).
The Storm Core API provides at-least-once semantics through an acking mechanism described here [1]. The Trident API builds on top of that to support exactly-once semantics by essentially doing a de-dupe [2].
I'm not sure exactly why they don't claim to support this, since Trident is build on top of Storm's Core API.
[1] https://storm.apache.org/documentation/Acking-framework-impl... [2] https://storm.apache.org/documentation/Trident-state.html
@filereaper:
Assuming you are referring to Spark streaming, forget about any benchmarks you may have seen. Either can be faster than the other depending on how you configure it, and what your use case is. See my presentation on the subject here [3]. With either, you can configure yourself into a corner and screw your performance.
Performance tuning distributed systems is a mysterious art. As is benchmarking. Unfortunately, that fact is frequently exploited for "benchmarketing" purposes. Don't trust any benchmark but your own unless it is fully open-sourced (including configuration).
[3] http://www.slideshare.net/ptgoetz/apache-storm-vs-spark-stre...
@vicaya
Version numbers don't necessarily equate to code quality, performance, or stability. I've seen many projects bump to 1.0 only for marketing purposes.
if you have the same requirement as twitter, nice -- but if you detour even a little bit, you are going to have a shitty time (for example, for thrift, i wanted to use buffered codec instead of framed and there was not a single document explaining how to do it. i spent ~2h perusing unit tests to find a way which i don't know if it's correct or not).
Disclaimer: I use both Mesos and Aurora.
http://www.digitaltrends.com/social-media/twitters-war-on-th...
But what about services that Twitter replicates? The latest applications victimized by Twitter’s in-house team are photo-sharing platforms like TwitPic and YFrog
Back in May, we first heard that Twitter planned to launch its own image posting tool. Up until that point, third party apps were giving Twitter that function. Now that Twitter’s own photo tool has launched, it’s virtually declaring war on the developers who were responsible for creating what has been a very popular aspect of the site.
But you're right, if you take time to develop an app using their API, it might just get their attention and then they'll develop their own version and then cut you completely out.
Twitter can absolutely block you from their API. Or come out with a competing product that integrates with their own systems better than your product does. But if you're using Storm (or Scalding or whatever), there's really nothing they can do to screw with that.
The two really can't be compared.