Skip to main content

CI status Latest PyPI release Supported Python versions MIT license

Introduction

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 pipelines.

Requirements & Installation

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 riko

riko installs a slim core by default. View the installation doc for advanced installation options.

Quick start

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, Sources, SyncPipe
>>>
>>> ### Set the pipe configurations ###
>>> #
>>> # Notes:
>>> #   1. look up cached html file in the `data` directory
>>> #   2. fetch text in the 'body' tag and strip html tags
>>> #   3. replace newlines with spaces and assign the result to 'content'
>>> #   4. split text in words using whitespace as the delimiter
>>> #   5. count the number of times each word 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(Sources.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}

Motivation

Why I built riko

I wanted a small-footprint, pure-Python library for processing data streams. In particular, I wanted to fetch RSS feeds and web pages and process records without needing to deploy a scheduler, cluster, or message queue.

The basic idea is deliberately simple: dictionary-like records flow through configurable pipes. Pipelines can run synchronously, asynchronous via async/await, or parallelized across threads or processes.

Why you should use riko

riko is a good fit when you want a batteries included, reusable, data-processing abstraction.

In particular, riko provides:

  • a pure-Python, embedded execution model with no required external services

  • a library of configuration-driven pipes for filtering, sorting, parsing, transforming, aggregating, and composing streams

  • first-class RSS/Atom and web-content processing

  • synchronous and asynchronous APIs

  • local thread and process-pool execution

  • lazy iterator-oriented processing

  • simple Python or JSON pipeline configuration and definition

  • tools to inspect, execute, and compile pipelines

Why you shouldn’t use riko

riko does not try to be a distributed stream-processing engine, durable workflow scheduler, or dataframe query engine.

It is usually not the right tool when you need:

  • execution across a cluster

  • durable keyed state and recovery after worker failure

  • persistent scheduling, retries, or task dependency management

  • a workflow service/UI

  • event-triggered infrastructure automation

  • dataframe-scale columnar analytics or query optimization

riko can instead run inside a worker or task managed by such systems.

Choosing riko

Several Python projects overlap with riko, but they optimize for different parts of the data-processing problem.

Project | Distinctive strength

Prefer it for…

dlt | Declarative, schema-aware
ingestion

moving data from REST APIs into warehouses/lakes/databases

Singer | Standardized taps and targets

replicating data from various sources into many destinations

Bytewax | Stateful streaming runtime

keyed state, recovery, workers, or distributed stream processing

Bonobo | Injectable services and I/O

traditional ETL graphs and runtime-injected infrastructure

Streamz | Continuous stream graphs

push-oriented streams, branching, backpressure, or live windows

petl | Rich lazy table algebra

joins, reshaping, and data-quality operations

riko | Config-driven pipelines

broad library of reusable, JSON serializable pipes

The closest comparison depends on what part of riko you care about.

dlt is a Python ingestion framework/library with similarities to riko in REST ingestion, incremental extraction, schema-aware loading, and Python-native data handling. dlt primarily allows you to “get data out a source reliably and into a well-structured destination.” It provides primitives for pagination, auth, and schema normalization. This contrasts with riko’s main use-case of processing and composing streams of records.

Singer is a connector protocol that standardizes how sources (taps) and destinations (targets) exchange records, schemas, and replication state. It overlaps with riko at the extraction and data-movement boundaries. Singer is a better fit when the primary goal is source-to-destination replication. riko instead places more emphasis on transforming and composing records.

Bytewax is the natural direction when a workload grows beyond riko’s intended scope and requires durable keyed state, recovery, or distributed stream processing.

petl and Bonobo, like riko, are both lightweight ETL libraries. Compared to riko, petl is more table-oriented and provides a deeper relational/data-wrangling vocabulary. Bonobo centers execution around an ETL graph of transformation nodes.

Streamz overlaps most with riko’s stream-composition and fan-out model, but places more emphasis on continuous push-based streams, windowing, and reactive dataflow.

riko provides more “batteries included” data-processing vocabulary. It exposes common operations (filtering, truncating, searching, etc.) as configurable, reusable pipes rather than requiring a Python callable. riko also provides first-class support for web-content (RSS/Atom feeds, HTML/XML, and JSON) and a simple JSON-based pipeline definition format.

Design Principles

Overview

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=...)

Core concepts

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, Transforms
>>>
>>> items = [{'title': 'riko pt. 1'}, {'title': 'riko pt. 2'}]
>>> stream = SyncPipe(Transforms.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, Transforms
>>>
>>> items = [{'title': 'riko pt. 1'}]
>>> stream = SyncPipe(Transforms.HASH, items, field='title')
>>> next(stream)['hash']
1104819838

Some processors, e.g., tokenizer, return multiple results.

>>> from riko import SyncPipe, Transforms
>>>
>>> items = [{'title': 'riko pt. 1'}]
>>> stream = SyncPipe(Transforms.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, Transforms
>>>
>>> items = [{'title': 'riko pt 1'}, {'title': 'riko pt 2'}]
>>> list(SyncPipe(Transforms.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, Transforms
>>>
>>> stream = SyncPipe(Transforms.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 Sources, SyncPipe
>>>
>>> attrs = {'key': 'title', 'value': 'riko pt. 1'}
>>> next(SyncPipe(Sources.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, Sources, Transforms
>>>
>>> metadata = get_module_metadata(Sources.FETCHPAGE)
>>> metadata.name, metadata.type, metadata.subtype
('fetchpage', 'processor', 'source')
>>> metadata = get_module_metadata(Transforms.COUNT)
>>> metadata.name, metadata.type, metadata.subtype
('count', 'operator', 'aggregator')

Note: type and subtype are mutually exclusive: a subtype implies its type.

SyncPipe/AsyncPipe perform this check for you to allow for convenient method chaining and transparent parallelization.

>>> from riko import Sources, SyncPipe
>>>
>>> attrs = [
...     {'key': 'title', 'value': 'riko pt. 1'},
...     {'key': 'content', 'value': "Let's talk about riko!"}
... ]
>>> flow = SyncPipe(Sources.ITEMBUILDER, conf={'attrs': attrs}).hash()
>>> item = next(flow)
>>> item['title'], item['content'], item['hash']
('riko pt. 1', "Let's talk about riko!", 197222720)

The | operator chains the same way. It takes a module name or a (name, conf) tuple. The later is handy when the next pipe’s name is computed. A name may be a plain string or a member of the typed discovery tree (Sources/Transforms/ Sinks).

>>> from riko import Sources, SyncPipe, Transforms
>>>
>>> attrs = [
...     {'key': 'title', 'value': 'riko pt. 1'},
...     {'key': 'content', 'value': "Let's talk about riko!"}
... ]
>>> conf = {'attrs': attrs}
>>> item = next(SyncPipe(Sources.ITEMBUILDER, conf=conf) | Transforms.HASH)
>>> item['title'], item['hash']
('riko pt. 1', 197222720)

View the Cookbook for advanced examples including how to wire in values from other pipes or accept user input.

Usage

riko can be used directly as a Python library.

Usage Index

Fetching data

riko can fetch data such as HTML, JSON, CSV, etc. from both local and remote filepaths via source pipes:

>>> from riko import get_path, Sources, SyncPipe
>>>
>>> stream = SyncPipe(Sources.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.

Synchronous processing

riko can modify a stream via transformer, composer, and aggregator pipes:

>>> from riko import get_path, Sources, 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(Sources.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.

Parallel processing

An example using riko’s parallel API to spawn a ThreadPool [1]

>>> from riko import get_path, Sources, 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(Sources.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

Asynchronous processing

To enable asynchronous processing, you must install the async extra.

python -m pip install "riko[async]"
>>> from riko import AsyncPipe, get_path, issync, run, Sources
>>>
>>> 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(Sources.FETCH, conf=fetch_conf)           # 1
...             .filter(conf={'rule': filter_rule}))            # 2
...
...     print(next(stream)['title'])
>>>
>>> print('Donations') if issync else run(main)
Donations

Built-in pipes

riko ships 52 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

feed and network-backed transformations

Sinks & writers

write

serialize a stream to a file (materializes it)

Pipeline lifecycle

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, Transforms
>>>
>>> flow = SyncPipe(Transforms.HASH, source=[{'content': 'a'}, {'content': 'b'}])
>>> flow.state
<PipeState.NEW: 'new'>
>>> len(list(flow))
2
>>> flow.state
<PipeState.EXHAUSTED: 'exhausted'>
>>> flow.exhausted
True

See the Cookbook for pool cleanup and the full state model.

Command-line Interface

riko provides a command, run-pipe, to execute pipelines. A pipeline 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 Sources, 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(Sources.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 pipeline 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 pipelines. Type run-pipe demo and you should see the following output:

Deadline to clear up health law eligibility near
682

Contributing

Please 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.

Credits

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

More Info

  • FAQ — the complete built-in pipe and file-format catalog

  • Cookbook — progressively organized, runnable recipes

  • DAG format — compact and full JSON pipeline formats

  • Migration guide — upgrading from the older versions or the legacy branch

  • Changelog — release notes

  • Contributing doc — contribution and issue-reporting guidance

  • issue tracker — bugs, feature proposals, and questions

Project Structure

┌── _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.lock

License

riko is distributed under the MIT License.

Download files

Download the file for your platform. If you're not sure which to choose, learn more about installing packages.

Source Distribution

riko-0.76.0.tar.gz (1.2 MB view details)

Uploaded Source

Built Distribution

If you're not sure about the file name format, learn more about wheel file names.

riko-0.76.0-py3-none-any.whl (1.3 MB view details)

Uploaded Python 3

File details

Details for the file riko-0.76.0.tar.gz.

File metadata

  • Download URL: riko-0.76.0.tar.gz
  • Upload date:
  • Size: 1.2 MB
  • Tags: Source
  • Uploaded using Trusted Publishing? Yes
  • Uploaded via: twine/7.0.0 CPython/3.13.14

File hashes

Hashes for riko-0.76.0.tar.gz
Algorithm Hash digest
SHA256 ad5b5816898d641f5edddde49ec9292c43abb641f1e13e03f52087edbb77bcf0
MD5 8f5add78d0aeda84941ee2463e7c6ac0
BLAKE2b-256 513262bb5972eb88720917fbc572dd5e9c80fa9b422b459d28e900b38bda9dd4

See more details on using hashes here.

Provenance

The following attestation bundles were made for riko-0.76.0.tar.gz:

Publisher: publish.yml on nerevu/riko

Attestations: Values shown here reflect the state when the release was signed and may no longer be current.

File details

Details for the file riko-0.76.0-py3-none-any.whl.

File metadata

  • Download URL: riko-0.76.0-py3-none-any.whl
  • Upload date:
  • Size: 1.3 MB
  • Tags: Python 3
  • Uploaded using Trusted Publishing? Yes
  • Uploaded via: twine/7.0.0 CPython/3.13.14

File hashes

Hashes for riko-0.76.0-py3-none-any.whl
Algorithm Hash digest
SHA256 d16e1ce37b681ead65406c5b039b94debefa9ce29b87801cb130d54c1b5bcbc3
MD5 c7c28525d1973293afdf59359766e5de
BLAKE2b-256 1989ad41f8a4c265769c6279cf20bb34659ca44e27d54e2de35f557604e0eaee

See more details on using hashes here.

Provenance

The following attestation bundles were made for riko-0.76.0-py3-none-any.whl:

Publisher: publish.yml on nerevu/riko

Attestations: Values shown here reflect the state when the release was signed and may no longer be current.

Release history Release notifications | RSS feed

0.76.1

2 files

This release

0.76.0 This release

2 files

0.74.1

2 files

0.74.0

2 files

0.73.1

2 files

0.73.0

2 files

0.72.2

2 files

0.71.2

2 files

0.71.1

2 files

0.71.0

2 files

0.70.0

2 files

0.69.0

2 files

0.68.1

2 files

0.67.0

2 files

0.66.0

2 files

0.65.0

2 files

0.64.3

2 files

0.64.2

2 files

0.64.1

2 files

0.64.0

2 files

0.63.0

2 files

0.62.2

2 files

0.62.1

2 files

0.62.0

2 files

0.61.4

2 files

0.61.2

2 files

0.61.1

2 files

0.60.4

2 files

0.60.3

2 files

0.60.2

2 files

0.60.1

2 files

0.60.0

2 files

0.59.1

2 files

0.59.0

2 files

0.58.0

2 files

0.57.0

2 files

0.56.3

2 files

0.56.2

2 files

0.56.1

2 files

0.56.0

2 files

0.55.0

2 files

0.54.1

2 files

0.54.0

2 files

0.53.0

2 files

0.52.3

2 files

0.52.2

1 file

0.52.1

2 files

0.51.0

2 files

0.50.0

2 files

0.49.2

2 files

0.47.0

2 files

0.46.1

2 files

0.46.0

2 files

0.45.1

2 files

0.45.0

2 files

0.44.0

2 files

0.43.1

2 files

0.43.0

2 files

0.42.0

2 files

0.41.0

2 files

0.40.1

2 files

0.39.0

2 files

0.38.0

2 files

0.37.0

2 files

0.36.0

2 files

0.35.3

2 files

0.35.1

2 files

0.35.0

2 files

0.33.0

2 files

0.32.1

2 files

0.30.0

2 files

0.29.0

2 files

Anthropic, PBC Visionary sponsor Bloomberg Visionary sponsor Hudson River Trading Visionary sponsor Meta Visionary sponsor NVIDIA Visionary sponsor Microsoft Sustainability sponsor Depot Continuous Integration AWS Cloud computing and Security Sponsor Datadog Monitoring Fastly CDN Google Download Analytics Sentry Error logging StatusPage Status page