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 :-(
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 :-(
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...
[1] https://github.com/nerevu/riko
Edit: HN submission
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...
https://www.influxdata.com/time-series-platform/kapacitor/
Disclaimer: I am the author.
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.
You can have a look here : https://hexdocs.pm/gen_stage/Experimental.Flow.html
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.