Scalable Stream Processing: A Survey of Storm, Samza, Spark and Flink
medium.com
medium.com
0: http://beam.incubator.apache.org/
1: http://softwareengineeringdaily.com/2016/08/19/apache-beam-w...
This sounds like a stupid comment, but when the underlying services are all evolving as quickly as they do in this space programming against an abstraction layer means you need to wait for (often sorely needed) new features.
(Frances Perry is an engineer on Beam at Google, so it would be surprising if she recommended against it)
You have a point that abstraction layers can limit you since it's another compatibility layer, but looking at what's done in actual practice I don't see it being that bad.
If anything, it should be evaluated on a case by case basis.
Flink is actually better tech overall for my use case, but a lot of customers want spark streaming since it's already installed. Having beam where we can do both is kinda nice.
We jumped from 1.5 to 1.6 because of algorithms in MLLib (although that turned out to be a bit of a disappointment).
Something I can point to is a modified Linear Road benchmark: https://github.com/IBMStreams/benchmarks/tree/master/Streams...
This benchmark was made at the request of a potential customer. Our implementation scaled to 200 "lanes." The other systems tested did not scale past 50. Unfortunately, that's as much as I feel I can say until I speak with some of the people involved in the comparison.
Development community: https://developer.ibm.com/streamsdev/
Some pointers to posts I've made in the development community focusing on the language and performance: http://www.scott-a-s.com/streams-posts/
Academic paper on the language; this is an IBM technical report, a version of this will be published in TOPLAS: http://hirzels.com/martin/papers/tr14-rc25486-spl.pdf
Brief academic paper on the systems aspects of the language: http://hirzels.com/martin/papers/debull15-spl.pdf
Something like spark or kafka streaming but that doesn't depend on hundreds of megabytes of java stuffs?
I just want to do basic windowing/counting against data streams.. I don't want to be part of a gigantic java ecosystem :-(
https://www.influxdata.com/time-series-platform/kapacitor/
Disclaimer: I am the author.
the esper stuff was on my radar, but something that works with existing postgresql tools is a huge plus.
Doesn't that get in the way of low latency and availability?
[1] http://enterprise.pipelinedb.com/docs/two-phase.html#two-pha...
[1] https://github.com/nerevu/riko
Edit: HN submission
I also updated the dashboard to match [2], including Riemann's query language [3].
I haven't had a chance to update the docs yet, but the tests have a ton of examples [4]. Reactors are just Transform streams now so it's possible to use any duplex/transform stream in the pipeline. Docs coming ASAP.
[1]: https://github.com/nextorigin/godot2
[2]: https://github.com/nextorigin/godot2-dash
[3]: https://github.com/nextorigin/riemann-query-parser
[4]: https://github.com/nextorigin/godot2/tree/nextorigin/test/re...
The product also comes with dashboards so you can do ad-hoc analytics and data visualization on your streams and windows.
It has support for stream replay built on top of Kafka-backed persistent streams. It has a commercial license, but we're pretty flexible with startups
It's a bit harder to dive into java though.
The thing I have a use case for stream processing is basically glorified wordcount - just over a realtime 7 day window. The data processing code is the easy part, the tooling/runtime is the hard part.
I'm not sure the size of the install should be the primary thing you are concerned about.
I'd prefer the kafka streams approach of deployment, but I suppose for just having it run on a single box this would work.
It doesn't solve the data storage problem, so if you need that you need a way to map the same files to the same path on every node (eg, network drive, RSync).
If you are just doing streaming it might not matter though.
The ability to run the data pipeline in realtime but also replay older data against new code using something like kafka would be a huge plus.
Really the main thing my code is missing is the sliding window implementation, so I can either port to spark streaming which has all that stuff built in, or just implement my own window code.
windowedWordCounts = pairs.reduceByKeyAndWindow(lambda x, y: x + y, lambda x, y: x - y, 30, 10)
Search for "window operations" on http://spark.apache.org/docs/latest/streaming-programming-gu.... Unless you meant something different?You can replay streams against Spark, too. streamingContext.textFileStream will stream data from files dumped in a directory - to replay them, just dump them there again.
reduceByKeyAndWindow is what my simple non-spark proof of concept is missing. otherwise the code is basically the same.
my non spark code is essentially
while True:
data = {}
for line in input:
rec = parse(line)
data = aggregate(data, rec)
data = filter(is_bad, data)
pprint(data)
the spark version is 99% the same code: lines.map(parse).reduceByKeyAndWindow(add, sub, 3600, 60).
filter(is_bad).pprint()
Figuring out how to do stateful processing is a little tricky, but updateStateByKey seems to do what i need.. I need to dedup the output per key for time period t. Though, some recommendations are just to use something like redis or memcached which would work.You can have a look here : https://hexdocs.pm/gen_stage/Experimental.Flow.html