riko is a pure Python library for building data-processing streams.
riko combines reusable, configuration-driven modular pipes with synchronous,
asynchronous, and parallel execution APIs. It is particularly useful for
processing RSS feeds, web content, text, and structured files.
riko also supplies a command-line interface for executing flows, i.e.,
stream processors aka workflows.
riko has been tested and is known to work on Python 3.12, 3.13, and 3.14.
Install the latest published release from PyPI:
python -m pip install rikoriko installs a slim core by default. View the installation doc for advanced
installation options.
The following example fetches a webpage, splits its text into words, and counts the number of times each word appears.
>>> from riko import get_path, SyncPipe
>>>
>>> ### Set the pipe configurations ###
>>> #
>>> # Notes:
>>> # 1. look up cached file in the `data` directory
>>> # 2. fetch the text contained inside the 'body' tag of a web page and strip
>>> # html tags
>>> # 3. replace newlines with spaces and assign the result to 'content'
>>> # 4. tokenize the resulting text using whitespace as the delimiter
>>> # 5. count the number of times each token appears
>>>
>>> url = get_path('users.jyu.fi.html') # 1
>>> fetch_conf = {'url': url, 'start': '<body>', 'end': '</body>', 'detag': True}
>>> replace_conf = {
... 'rule': [{'find': '\r\n', 'replace': ' '}, {'find': '\n', 'replace': ' '}]
... }
>>>
>>> flow = (
... SyncPipe('fetchpage', conf=fetch_conf) # 2
... .strreplace(conf=replace_conf, assign='content') # 3
... .tokenizer(conf={'delimiter': ' '}, emit=True) # 4
... .count(conf={'count_key': 'content'}) # 5
... )
>>>
>>> next(flow)
{'Tidy': 1}
>>> next(flow)
{'your': 1}I wanted a small-footprint, pure-Python library for processing data streams. I wanted to fetch RSS feeds and perform custom transformations without the complexities of distributed compute engines, workflow schedulers, clusters, or message queues.
riko, a primarily pull-based in-process library, is the result.
riko provides a number of benefits / differences from other stream processing
applications:
- a small footprint (CPU and memory usage)
- native RSS/Atom support
- simple installation and usage
- a pure Python library supporting v3.12+
- built-in modular
pipesto filter, sort, and modifystreams
riko is usually not the right tool when you need: distributed execution,
durable scheduling, automatic retries, continual data monitoring, a workflow UI,
event-triggered actions, or query optimization.
The projects below overlap with riko in different ways. Some are embedded libraries, while others are distributed engines, workflow orchestrators, or data-integration platforms. RSS/Atom role distinguishes first-party feed support from functionality that requires a custom source, connector, tap, or task.
| Project | Primary model | Deployment | RSS/Atom | Best fit |
|---|---|---|---|---|
| riko | Python pipelines | Embedded, local process | Built in | Lightweight data and feed processing |
| RxPY | Reactive observables | Embedded library | Custom adapter | Push-based application events |
| Huginn | Persistent agents | Self-hosted app | Built in | UI-driven monitoring and automation |
| Apache Beam | Portable pipelines | Runner-dependent | Custom I/O | Portable batch and stream processing |
| Flink | Stateful streams | Distributed engine | Custom connector | Low-latency, stateful processing |
| Storm | Event topologies | Distributed engine | Custom spout | Low-latency event processing |
| Spark | DataFrame streams | Distributed engine | Custom connector | Large-scale streaming analytics |
| Luigi / Prefect | Task workflows | Scheduler and workers | External task | Scheduling, retries, and dependencies |
| Airbyte / Meltano | ELT platforms | Connectors or plugins | Connector/tap dependent | Data integration and repeatable ELT pipelines |
Choose riko when a pipeline should run directly inside a Python application without a separate scheduler, service, or cluster. It provides first-party RSS/Atom processing and supports synchronous, asynchronous, and thread-pooled local execution.
Choose RxPY for reactive event composition, Huginn for persistent
UI-managed automation, and Flink, Storm, Spark, or Apache Beam when
distributed execution is required. Luigi and Prefect orchestrate tasks, while Airbyte
and Meltano focus on data-integration workflows. Meltano commonly runs Singer taps and
targets and can add scheduling through an orchestration plugin. These tools may
run riko as one step in a larger workflow rather than replace its in-process
transformation API.
Here's the riko vocabulary at a glance:
| Term | Meaning | Example |
|---|---|---|
item |
one dictionary-like record | {'title': 'Example'} |
stream |
an iterator of item |
iter([{'title': 'Example'}]) or SyncPipe |
pipe |
a configured stream operation | join, slugify, uniq |
operator |
a pipe that consumes a stream |
count, filter, reverse |
processor |
a pipe that consumes an item |
urlparse, fetch, hash |
splitter |
a pipe returning multiple streams |
split |
flow / pipeline |
a chain of configured pipes |
SyncPipe(...).count() |
Context |
runtime inputs + ExecutionMode |
Context(inputs=...) |
The primary data structures in riko are the item and stream. An item
is just a Python dictionary, and a stream is an iterator of item. You can
create a stream manually with something as simple as
iter([{'content': 'hello world'}]). You manipulate streams in
riko via pipes. A pipe is simply a function that accepts either a
stream or item, and returns a stream.
Through SyncPipe and AsyncPipe classes, pipes are composable: the output of
each pipe is the input to the next pipe.
riko pipes come in three types: processor, operator, and splitter.
An operator operates on a stream and is unable to handle individual items.
E.g., count, filter, and reverse.
>>> from riko import SyncPipe
>>>
>>> items = [{'title': 'riko pt. 1'}, {'title': 'riko pt. 2'}]
>>> stream = SyncPipe('reverse', items)
>>> next(stream)
{'title': 'riko pt. 2'}A processor processes an individual item and can be parallelized across
threads or processes. E.g., fetchsitefeed, hash, itembuilder, and regex.
>>> from riko import SyncPipe
>>>
>>> items = [{'title': 'riko pt. 1'}]
>>> stream = SyncPipe('hash', items, field='title')
>>> next(stream)['hash']
1104819838Some processors, e.g., tokenizer, return multiple results.
>>> from riko import SyncPipe
>>>
>>> items = [{'title': 'riko pt. 1'}]
>>> stream = SyncPipe('tokenizer', items, conf={'delimiter': ' '}, field='title')
>>> list(stream)
[{'content': 'riko'}, {'content': 'pt.'}, {'content': '1'}]operators are split into sub-types: aggregator
and composer. aggregators, e.g., count, combine
all items of an input stream into a new stream with a single item;
while composers, e.g., filter, create a new stream containing
some or all items of an input stream.
>>> from riko import SyncPipe
>>>
>>> items = [{'title': 'riko pt 1'}, {'title': 'riko pt 2'}]
>>> list(SyncPipe('count', items))
[{'count': 2}]Astute observers may have noticed from the "Word Count" example up top, that count
can return multiple items if you pass in the count_key config option.
>>> from riko import SyncPipe
>>>
>>> stream = SyncPipe('count', items, conf={'count_key': 'title'})
>>> list(stream)
[{'riko pt 1': 1}, {'riko pt 2': 1}]processors are parallelizable and split into sub-types of source and
transformer. A source, e.g., itembuilder, can create a stream, while
a transformer, e.g. hash can only transform a source item.
>>> from riko import SyncPipe
>>>
>>> attrs = {'key': 'title', 'value': 'riko pt. 1'}
>>> next(SyncPipe('itembuilder', conf={'attrs': attrs}))
{'title': 'riko pt. 1'}The following table summarizes these observations:
| Type | Sub-type | Meaning | Example |
|---|---|---|---|
| processor | source |
creates a stream |
itembuilder, fetch |
transformer |
manipulates an item |
hash, rename, regex |
|
| operator | composer |
selects/orders a stream |
filter, sort, union |
aggregator |
summarizes a stream |
count, sum |
|
| splitter | splitter |
copies a stream |
split |
Note: Since some pipes support more than one subtype depending on their options,
view the FAQ for steps on runtime discovery via discovering modules
If you are unsure of the type of pipe you have, check its metadata.
>>> from riko import get_module_metadata
>>>
>>> metadata = get_module_metadata('fetchpage')
>>> metadata.name, metadata.type, metadata.subtype
('fetchpage', 'processor', 'source')
>>> metadata = get_module_metadata('count')
>>> metadata.name, metadata.type, metadata.subtype
('count', 'operator', 'aggregator')SyncPipe/AsyncPipe perform this check for you to allow for convenient method
chaining and transparent parallelization.
>>> from riko import SyncPipe
>>>
>>> attrs = [
... {'key': 'title', 'value': 'riko pt. 1'},
... {'key': 'content', 'value': "Let's talk about riko!"}
... ]
>>> flow = SyncPipe('itembuilder', conf={'attrs': attrs}).hash()
>>> item = next(flow)
>>> item['title'], item['content'], item['hash']
('riko pt. 1', "Let's talk about riko!", 197222720)View the Cookbook for advanced examples including how to wire in values from other pipes or accept user input.
Note: type and subtype are mutually exclusive: a subtype implies its type.
riko can be used directly as a Python library.
- Fetching data
- Synchronous processing
- Parallel processing
- Asynchronous processing
- Built-in pipes
- Pipeline lifecycle
riko can fetch data such as HTML, JSON, CSV, etc. from both local and remote
filepaths via source pipes:
>>> from riko import get_path, SyncPipe
>>>
>>> stream = SyncPipe('fetch', conf={'url': get_path('feed.xml')})
>>> item = next(stream)
>>> {'author', 'content', 'id', 'link', 'published', 'summary', 'title'} <= set(item)
True
>>> item['title'], item['author'], item['id']
('Donations', {'name': 'WriteToReply', 'uri': None}, 'http://writetoreply.org/?page_id=111')View the FAQ for a complete list of supported file types and protocols; and Fetching data and feeds for more examples.
riko can modify a stream via transformer, composer, and aggregator
pipes:
>>> from riko import get_path, SyncPipe
>>>
>>> fetch_conf = {'url': get_path('feed.xml')}
>>> filter_rule = {'field': 'title', 'op': 'contains', 'value': 'a'}
>>>
>>> # The following flow will:
>>> # 1. fetch a (cached) RSS feed
>>> # 2. filter for items with an 'a' in the title
>>> # 3. sort the items ascending by title
>>> #
>>> # Note: sorting is not lazy so take caution when using this pipe
>>>
>>> flow = (
... SyncPipe('fetch', conf=fetch_conf) # 1
... .filter(conf={'rule': filter_rule}) # 2
... .sort(conf={'rule': {'field': 'title'}}) # 3
... )
>>>
>>> next(flow)['title']
'Donations'View pipes for a complete list of available pipes.
An example using riko's parallel API to spawn a ThreadPool [1]
>>> from riko import get_path, SyncPipe
>>>
>>> fetch_conf = {'url': get_path('feed.xml')}
>>> filter_rule = {'field': 'title', 'op': 'contains', 'value': 'a'}
>>>
>>> # The following flow will:
>>> # 1. fetch a (cached) RSS feed
>>> # 2. filter for items with an 'a' in the title, in parallel (4 workers)
>>> #
>>> # Note: no point in sorting after the filter since parallel fetching doesn't
>>> # guarantee order
>>> flow = (
... SyncPipe('fetch', conf=fetch_conf, parallel=True, workers=4) # 1
... .filter(conf={'rule': filter_rule}) # 2
... )
>>>
>>> sorted(item['title'] for item in flow)[:3]
['Donations', 'FAQ', 'General Comments']Notes
| [1] | You can instead enable a ProcessPool by additionally passing threads=False to SyncPipe, i.e., SyncPipe('fetch', conf={'url': url}, parallel=True, threads=False). |
To enable asynchronous processing, you must install the async extra.
python -m pip install "riko[async]">>> from riko import AsyncPipe, get_path, issync, run
>>>
>>> fetch_conf = {'url': get_path('feed.xml')}
>>> filter_rule = {'field': 'title', 'op': 'contains', 'value': 'a'}
>>>
>>> # The following flow will:
>>> # 1. fetch a (cached) RSS feed
>>> # 2. filter for items with an 'a' in the title
>>>
>>> async def main():
... stream = await (
... AsyncPipe('fetch', conf=fetch_conf) # 1
... .filter(conf={'rule': filter_rule})) # 2
...
... print(next(stream)['title'])
>>>
>>> print('Donations') if issync else run(main)
Donationsriko ships 51 built-in pipes. The table below summarizes them.
| Group | Representative pipes | Purpose |
|---|---|---|
| Sources & readers | itembuilder, fetch, fetchtable, csv |
build items from config, feeds, files, or input |
| Selection & ordering | filter, sort, truncate, uniq |
select, order, dedupe, or bound a stream |
| Text & field transforms | regex, rename, strreplace, tokenizer |
extract and transform string / item fields |
| Type & numeric transforms | typecast, simplemath, dateformat, hash |
convert types and derive fields |
| Aggregation & combination | count, sum, join, union, split |
summarize, merge, join, or copy streams |
| Control & extension | loop, udf, send, receive |
run submodules, call funcs, fan out items |
| Feed & location helpers | fetchsitefeed, exchangerate, geolocate |
feeds and network-backed transformations |
SyncPipe/AsyncPipe represent a single execution: iterating one
consumes the stream, and iterating it again yields an empty stream. Read the
state/exhausted/closed/failed properties to inspect a pipe. Use it as a
context manager (or call close()/terminate()) to release a parallel pipe's
worker pool deterministically.
>>> from riko import SyncPipe
>>>
>>> flow = SyncPipe('hash', source=[{'content': 'a'}, {'content': 'b'}])
>>> flow.state
<PipeState.NEW: 'new'>
>>> len(list(flow))
2
>>> flow.state
<PipeState.EXHAUSTED: 'exhausted'>
>>> flow.exhausted
TrueSee the Cookbook for pool cleanup and the full state model.
riko provides a command, run-pipe, to execute workflows. A
workflow is simply a file containing a function named pipe that creates
a flow and processes the resulting stream. E.g., flow.py
from riko import SyncPipe
conf1 = {'attrs': [{'value': 'https://google.com', 'key': 'content'}]}
conf2 = {'rule': [{'find': 'com', 'replace': 'co.uk'}]}
def pipe(test=False):
kwargs = {'conf': conf1, 'test': test}
flow = SyncPipe('itembuilder', **kwargs).strreplace(conf=conf2)
for i in flow:
print(i)CLI Usage
usage: run-pipe [pipeid] [-p PATH]
description: Runs a riko pipe
- positional arguments:
- pipeid The workflow to run from the examples directory.
- optional arguments:
-h, --help show this help message and exit -p, --path PATH Path to a pipe file to run, e.g. flow.py. -a, --async Load async pipe. -t, --test Run in test mode (uses default inputs).
Now to execute flow.py, type the command run-pipe --path flow.py. You should
then see the following output in your terminal:
{'content': 'https://google.com', 'strreplace': 'https://google.co.uk'}run-pipe will also search the examples directory for workflows. Type
run-pipe demo and you should see the following output:
Deadline to clear up health law eligibility near
682Please mimic the coding style/conventions used in this repo. If you add new classes or functions, please add the appropriate docstrings with examples. Also, make sure the linter and tests pass.
View Contributing doc for more details.
riko started out as a fork of pipe2py which translated a Yahoo! Pipe [#] into
python code. riko has since diverged so much from pipe2py that little of the
original code-base remains.
Notes
| [2] | Discontinued in 2015, Yahoo! Pipes was a user friendly web application used to aggregate, manipulate, and mashup content from around the web. You can view what remains |
- FAQ — the complete built-in
pipeand file-format catalog - Cookbook — progressively organized, runnable recipes
- DAG format — compact and full JSON
workflowformats - Migration guide — upgrading from the older versions or the
legacybranch - Changelog — release notes
- Contributing doc — contribution and issue-reporting guidance
- issue tracker — bugs, feature proposals, and questions
┌── _docs/* (internal documentation)
├── docs
│ ├── AUTHORS.rst
│ ├── CHANGES.rst
│ ├── COOKBOOK.rst
│ ├── DAG_FORMAT.rst
│ ├── FAQ.rst
│ ├── INSTALLATION.rst
│ ├── MIGRATION.rst
│ └── ROADMAP.md
├── examples/*
├── riko
│ ├── __init__.py (stable public API)
│ ├── api.py (stable API re-export hub)
│ ├── autorss.py, cast.py, currencies.py, dates.py, locations.py, pprint2.py, topsort.py
│ ├── collections.py (SyncPipe, AsyncPipe, SyncCollection, AsyncCollection)
│ ├── compile.py (JSON pipe → executable pipeline / Python module)
│ ├── context.py (Context, ExecutionMode)
│ ├── dotdict.py
│ ├── paths.py (get_path / get_abspath)
│ ├── parsers.py (sync XML/HTML parsing)
│ │
│ ├── _*.py (private helpers: _feed, _io, _iterutils, _objectify,
│ │ _serialize, _strutils, _logging)
│ ├── _pubsub/ (sync + async pub/sub hubs backing send/receive)
│ ├── bado/ (async backend: __init__, io, itertools, mock, _util)
│ ├── cli/ (manage, run-pipe, benchmark, compile, convert-dag, gen-config)
│ ├── data/*
│ ├── ext/ (extension API: decorators, protocols)
│ ├── modules/* (the built-in pipes)
│ ├── templates/* (codegen Jinja templates)
│ └── types/ (compile, general, modules, values, configs, guards)
├── tests
│ ├── __init__.py
│ ├── conftest.py
│ ├── dags/* (bare-bones DAG fixtures)
│ ├── functional/*
│ ├── internal/*
│ ├── pipelines/* (JSON pipe definitions)
│ ├── public/*
│ └── pypipelines/* (expected generated Python modules)
├── CLAUDE.md
├── conftest.py
├── CONTRIBUTING.rst
├── LICENSE
├── pyproject.toml
├── README.rst
└── uv.lockriko is distributed under the MIT License.