Show HN: Riko – A Python stream processing engine modeled after Yahoo! Pipes
github.com
github.com
Out of the box, `riko` can read csv/xml/json/html files; create text and data based flows via modular pipes; parse and extract RSS/ATOM feeds; and bunch of other neat things. You can think of `riko` as a poor man's Spark/Storm... stream processing made easy!
Feedback welcome so let me know what you think!
Resources: FAQ [3], cookbook [4], and ipython notebook [5]
Quickie Demo:
>>> from riko.modules import fetch
>>>
>>> stream = fetch.pipe(conf={'url': 'https://news.ycombinator.com/rss'})
>>> item = next(stream)
>>> item['title'], item['link']
('Master Plan, Part Deux', 'https://www.tesla.com/blog/master-plan-part-deux')
[1] https://web.archive.org/web/20150930021241/http://pipes.yaho...[2] https://github.com/ggaughan/pipe2py/
[3] https://github.com/nerevu/riko/blob/master/docs/FAQ.rst
[4] https://github.com/nerevu/riko/blob/master/docs/COOKBOOK.rst
[5] http://nbviewer.jupyter.org/github/nerevu/riko/blob/master/e...
Yahoo Pipes was a nice project, but as its popularity grew, it started getting blocked more and more. It was also hard to build and maintain pipelines with more than a few steps.
[1] https://github.com/olviko/RssPercolator
[2] https://github.com/olviko/RssPercolator/blob/master/RssPerco...
But since generators can only be "pulled" into one destination, you have to copy a stream (subsequently converting it into a list) if you want more than one destination [5]. This works fine if the data can fit in memory, but if it can't then you're out of luck!
[1] http://www.dabeaz.com/coroutines/copipe.py
[2] http://www.dabeaz.com/coroutines/
[3] http://www.dabeaz.com/generators/retuple.py
[4] http://www.dabeaz.com/generators
[5] https://github.com/nerevu/riko/blob/master/riko/modules/spli...
But, I see what you mean. I had to deal with similar issues in commercial projects and the "pull" model (generators in Python ~ "yield return" in C#) almost never a good idea, especially when you have to have concurrent consumers. While callbacks are hard to combine, in C# it can be nicely abstracted with “async/await”, not sure how it is handled in Python, I stopped using it around 2.5
I've been working on a similar project and I've also found the push model easier.
1. The code is tiny with 90% of it dealing with RSS parsing and filtering. Using RX.NET wouldn't really simplify anything.
2. I wanted a library that I can integrate into my apps and run locally to avoid throttling, robots.txt and other BS Yahoo Pipes was suffering from.
I am also not a huge fan of RX... to put it mildly
Mind sharing why you don't like it ?
I personally like how it allows me to express complex high level operations cleanly. For example - I have a observable configuration variable that can come from different sources and the source change dynamically. I need to listen to latest source until a new one becomes active - in Rx I only need to push the new source trough IObservable<IObservable<ConfigurationValue>> and then use http://reactivex.io/documentation/operators/switch.html which returns IObservable<ConfigurationValue> which will push values from the latest source - Rx will handle unsubscribing from previous active source, synchronizing state and making sure everything is thread safe. And there are a bunch of operators like this that would be tedious and hard implement correctly with all the edge cases in a thread safe way - and here they are abstracted in to high level operators.
pretty much.... just without a GUI. My inspiration was method chaining [1] and the first implementation was this [2].
[1] http://martinfowler.com/articles/collection-pipeline/
[2] http://stackoverflow.com/questions/12172934/method-chaining-...
StreamCommandr essentially relies on the same principles and changes only external API and some internal processing. The idea is that any data table is like a stream of records so that we add new records and delete outdated records. Simultaniously, we evaluate other columns and tables by performaing potentially complex computations which are difficult to do in a record flow.
A more modest implementation would be Unix pipes, I think, where the data flows untyped.
Thanks!
It can handle parallel and distributed parts for you.
I personally prefer the functional approach much better. And if you compare the word count examples on the respective readmes [1, 2], you will see riko is much more succinct. But I suppose the verbosity of the other libraries come with benefits like scaling across a cluster of servers.
We use somewhat different concepts. I tend to think of streams as infinite, so it didn't occur to me to include something like a reverse pipe operator.
I'm a bit surprised, why is filter an operator rather than a processor? I would think filters usually apply per-item, not to a whole stream?
I havn't worked on it very much but I'm heading towards push-based, using 0MQ for distribution/parallel processing, and using asyncio, mostly because it plays nicely with 0MQ.
We are in agreement. reverse has a notice that it isn't lazy [1]. I prefer to include pipes that aren't lazy since it can be helpful in some cases (plus the goal is to include all pipes originally in Yahoo! Pipes). The vast majority of pipes work just fine on infinite streams [2].
> I'm a bit surprised, why is filter an operator rather than a processor? I would think filters usually apply per-item, not to a whole stream?
Just an implementation detail [3]. I agree it would be better if it were a processor since it could be parallelized. PRs welcome :).
> I haven't worked on it very much but I'm heading towards push-based, using 0MQ for distribution/parallel processing, and using asyncio, mostly because it plays nicely with 0MQ.
See my previous comments related to this [4, 5].
[1] https://github.com/nerevu/riko/blob/master/riko/modules/reve...
[2] Assuming you're not using the async or parallel mode
[3] https://github.com/nerevu/riko/blob/master/riko/modules/filt...
riko doesn't have a scheduler (although the original pipe2py has a json based one). However, I do plan to integrate with Airflow/Oozie/Luigi [1-3] in the future to make it easier to design workflows.
The notification system reminds me of Huggin [4]. Since riko is twisted based, it should be fairly straightforward to implement something similar for IRC/IMAP/FTP/etc.
[1] https://github.com/apache/incubator-airflow
As for riko more specifically, Beam will have soon a python sdk, but I'm unsure if there will be a python standalone runner. Maybe this is something to look into...
[1] https://www.oreilly.com/ideas/future-proof-and-scale-proof-y...
Just gave it a look. Took a while to find some examples with code, but once I did it made a bit more sense.
> Beam is runner-independent and you can take the same code and run it at scale on a cluster, wether it's spark, flink, or google cloud.
I thought that was pretty cool.
> As for riko more specifically, Beam will have soon a python sdk, but I'm unsure if there will be a python standalone runner. Maybe this is something to look into...
A python standalone runner would be very useful. Otherwise I'm hesitant to go much further since my goal is to have a pure python solution for working with streaming data. Most libraries require installing java and that is what I'd like to avoid.
[1] https://azkaban.github.io/ [2] https://developers.google.com/blockly/ [3] http://nodered.org/
Which interface do you think is more newbie friendly? My gut says blocks (maybe something a bit more simple/refined than blocky) are easier to grok, while wires allow for designing more complex workflows.
maybe figure out a few common workflows that people would make in riko or node-red, and mock up how they'd work/look in blocks vs. wiring/pipes.
Good idea, what are your workflows?
What kind of a demand is there for a pipes-kind of product or even a customizable/searchable rss/feed integrator?
How much would a typical user be willing to pay for it?
[1] http://mashable.com/2009/10/08/top-mashups/#0XwtqVCCXPq2
Really cool API, you should port this to concord! =)
i'd say major diff is dynamic topology. So during the pipeline execution you can add/remove workers for any stage.
Also each stage/operator can be written in any programming language.
Storm/Flink/SparkStreaming/etc... all have much higher level API's. We built the execution engine first, these great things (DSL, etc) should come soon. For example this API would be easy to support to execute on top (the pipe abstraction that is)
Here is an example of a DSL we prototyped in a couple hours.
Next up I think would be supporting custom sources/sinks such as Twitter, HDFS, RDMS, etc. What exactly would be involved in "porting" riko to concord and what would the advantages be for doing so?
Advantages are:
1. mesos integration, with that comes containerization support, multi tenancy, QoS, proper pipeline supervision, etc 2. Scheduling of pipelines. i.e.: Schedule them on 100 computers. 3. the outputs of your DAG could be consumed by other systems immediately and even written in different programming languages. So your python DSL could be the source to the Scala DSL at some point. so language interop 4. Available KV storage 5. Tracing (Zipkin) - ala Google Dapper. 6. Fast networking - C++ backed runtime is 1 order of magnitude faster than the python one.
What would be involved, is not much from what I can tell.
Each concord 'operator' is like a networked function.
so given a DAG, you could generate many operators internally, or literally write them to a file, i.e.: operator_one.py etc. The code generation or internal scheduling would be the glue that's needed.
if you ever become interested, ping me! would love to collab alex@concord.io
That's pretty impressive! So something like react/elm hot reloading?