Skip to main content

Lightweight async data pipelines using a publisher/subscriber pattern

Project description

io-chains

A lightweight Python library for building async data pipelines using a publisher/subscriber pattern. Chain together data sources, transformations, and consumers with minimal boilerplate.

Installation

pip install oj-io-chains

Core Concepts

A pipeline is built from Processors coordinated by a Chain:

[source] → Processor → Processor → Collector
                     ↘ Processor  (fan-out)

Each Processor:

  • Pulls items from a source (list, iterable, async generator, or callable)
  • Optionally transforms each item (sync or async)
  • Publishes to one or more subscribers

A Chain wires multiple Processors together and runs them concurrently — the caller just await chain().

For fan-in enrichment (joining data from multiple concurrent sources), use Enricher with Relation declarations and channel-tagged subscriptions.

To persist items to disk as they flow through the pipeline, use PersistenceLink.


Public API

from io_chains import Chain, Collector, Enricher, PersistenceLink, Processor, Relation, Skip, LinkMetrics

Classes

Processor

The core processing unit: source → transform → publish.

Processor(
    source=...,        # list, iterable, async gen, or callable (optional)
    processor=...,     # transform function — sync or async (optional)
    subscribers=[...], # downstream subscribers (optional)
    workers=1,         # number of concurrent workers
    batch_size=1,      # items per worker call
    on_error=None,     # callable(exc, item) — handle errors without stopping the stream
    on_metrics=None,   # callable(LinkMetrics) — called once on completion
    name="",           # label for metrics and logs
    queue_size=0,      # internal queue depth (0 = unbounded)
)

processor= return values

Return value Effect
Any value Published downstream
Skip Item dropped silently
Generator Expanded — each yielded value published
Async generator Expanded asynchronously
None (or omitted processor=) Original item passed through unchanged

Source types

Processor(source=[1, 2, 3])                      # list
Processor(source=(x * 2 for x in range(10)))     # generator expression
Processor(source=my_async_gen_func)              # async generator function (called automatically)
Processor(source=lambda: [1, 2, 3])              # callable returning iterable

Chain

Wires Processors (and sub-Chains) together and runs them concurrently. The caller just await chain().

chain = Chain(
    source=...,        # attached to the first link (optional)
    links=[...],       # ordered list of Processors or Chains
    subscribers=[...], # attached to the last link's output (optional)
)
await chain()

A Chain is itself a Linkable — it can be nested inside another Chain or used as a subscriber of an external Processor.


Collector

Buffers pipeline output for iteration after the pipeline completes.

results = Collector()

# async iteration (preferred — works inside a running event loop)
async for item in results:
    print(item)

# sync iteration (use after pipeline has completed)
for item in results:
    print(item)

PersistenceLink

A mid-chain tap: writes each item to an AsyncAbstractStore (from oj-persistence), then passes the item through unchanged.

The store's async context manager is entered automatically at run start and exited on completion, guaranteeing any buffered writes are flushed.

from oj_persistence.store.async_ndjson_file import AsyncNdjsonFileStore
from io_chains import PersistenceLink

PersistenceLink(
    store=AsyncNdjsonFileStore("data/output.ndjson"),
    key_fn=lambda item: str(item["id"]),  # extracts the store key from each item
    operation="upsert",   # "upsert" (default) | "create" | "update"
    on_error=None,        # callable(exc, item) — handle store errors without stopping the stream
)

Enricher and Relation

Fan-in join: collects items from multiple named channels, then streams primary items enriched via Relation declarations.

from io_chains import Enricher, Relation, Processor, Collector
from asyncio import create_task, gather

results = Collector()

enricher = Enricher(
    relations=[
        Relation(
            from_field="location_id",   # FK on the primary item
            to_channel="locations",     # channel holding related items
            to_field="id",              # field to match against
            attach_as="location",       # key added to enriched item
        ),
        Relation(
            from_field="episode_ids",
            to_channel="episodes",
            to_field="id",
            attach_as="episodes",
            many=True,                  # one-to-many: attach a list
        ),
    ],
    primary_channel="chars",
    subscribers=[results],
)

chars_link = Processor(source=chars_source)
locs_link = Processor(source=locations_source)
eps_link = Processor(source=episodes_source)

chars_link.subscribe(enricher, channel="chars")
locs_link.subscribe(enricher, channel="locations")
eps_link.subscribe(enricher, channel="episodes")

await gather(
    create_task(chars_link()),
    create_task(locs_link()),
    create_task(eps_link()),
    create_task(enricher()),
)

Relation parameters

Parameter Type Description
from_field str Field on the primary item whose value is the join key. For many=True, the value should be a list of keys.
to_channel str Channel name holding the related items.
to_field str Field on related items to match against from_field.
attach_as str Key added to the enriched primary item.
many bool True → one-to-many (attach list); False → one-to-one (attach single or None).
key_transform callable or None Optional transform applied to each key before lookup.

Skip

Sentinel returned by a processor= function to drop an item silently.

from io_chains import Skip

Processor(processor=lambda x: x if x > 0 else Skip)

LinkMetrics

Emitted once per Processor / PersistenceLink on completion via the on_metrics= callback.

@dataclass
class LinkMetrics:
    name: str
    items_in: int
    items_out: int
    items_skipped: int
    items_errored: int
    elapsed_seconds: float

Usage Examples

Simple transformation

import asyncio
from io_chains import Chain, Collector, Processor

async def main():
    results = Collector()
    await Chain(
        source=[1, 2, 3],
        links=[Processor(processor=lambda x: x * 2)],
        subscribers=[results],
    )()
    print([item async for item in results])  # [2, 4, 6]

asyncio.run(main())

Multi-stage pipeline

results = Collector()
await Chain(
    source=["a", "b", "c"],
    links=[
        Processor(processor=str.upper),
        Processor(processor=lambda x: f"item: {x}"),
    ],
    subscribers=[results],
)()
# ["item: A", "item: B", "item: C"]

Persist items to disk mid-pipeline

from oj_persistence.store.async_ndjson_file import AsyncNdjsonFileStore
from io_chains import Chain, Collector, PersistenceLink, Processor

results = Collector()
await Chain(
    source=fetch_records,
    links=[
        Processor(processor=normalize),
        PersistenceLink(
            store=AsyncNdjsonFileStore("data/records.ndjson"),
            key_fn=lambda r: str(r["id"]),
        ),
    ],
    subscribers=[results],
)()

Fan-out to multiple subscribers

sink1, sink2 = Collector(), Collector()
await Chain(
    source=[1, 2, 3],
    links=[Processor(processor=lambda x: x * 2)],
    subscribers=[sink1, sink2],
)()
# both sinks contain [2, 4, 6]

Nested Chains

normalise = Chain(links=[
    Processor(processor=lambda x: abs(x)),
    Processor(processor=lambda x: round(x, 2)),
])
stringify = Chain(links=[
    Processor(processor=lambda x: x * 100),
    Processor(processor=lambda x: f"{x:.0f}%"),
])
await Chain(source=[-0.156, 0.999, -0.301], links=[normalise, stringify], subscribers=[results])()
# ["16%", "100%", "30%"]

Observability

def log_metrics(m):
    print(f"{m.name}: {m.items_in} in, {m.items_out} out, {m.elapsed_seconds:.3f}s")

Processor(
    source=records,
    processor=transform,
    name="transform-stage",
    on_metrics=log_metrics,
)

Architecture

Subscriber (ABC)
└── Collector

Publisher (ABC)
└── Linkable(Publisher, Subscriber)  (ABC)
    ├── Link                          (internal base: queue, EOS, metrics)
    │   ├── Processor                 (source → transform → publish)
    │   └── PersistenceLink           (tap: write to store → pass through)
    ├── Chain                         (orchestrator: wires and runs Links)
    └── Enricher                      (fan-in: join multiple channels)

Development

# Install in editable mode
pip install -e .

# Run unit tests
python -m pytest test/unit -v

# Run user acceptance tests (requires network)
python -m pytest test/ua -v

# Lint
python -m ruff check io_chains/ test/

# Format
python -m ruff format io_chains/ test/

Project details


Download files

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

Source Distribution

oj_io_chains-0.0.4.tar.gz (15.9 kB view details)

Uploaded Source

Built Distribution

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

oj_io_chains-0.0.4-py3-none-any.whl (18.2 kB view details)

Uploaded Python 3

File details

Details for the file oj_io_chains-0.0.4.tar.gz.

File metadata

  • Download URL: oj_io_chains-0.0.4.tar.gz
  • Upload date:
  • Size: 15.9 kB
  • Tags: Source
  • Uploaded using Trusted Publishing? Yes
  • Uploaded via: twine/6.1.0 CPython/3.13.12

File hashes

Hashes for oj_io_chains-0.0.4.tar.gz
Algorithm Hash digest
SHA256 ee0916e5ae8326697a55781a03760c54cbfb81cccfc8d0e00e2ae0c43ea3a2d1
MD5 0bbef84405c6c9e0d9c59dad97401883
BLAKE2b-256 59b6ef38d2245f7e9122fdea6962ad784221d31c062df4205a29e4d252a4c589

See more details on using hashes here.

Provenance

The following attestation bundles were made for oj_io_chains-0.0.4.tar.gz:

Publisher: publish.yml on ownjoo-org/io-chains

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

File details

Details for the file oj_io_chains-0.0.4-py3-none-any.whl.

File metadata

  • Download URL: oj_io_chains-0.0.4-py3-none-any.whl
  • Upload date:
  • Size: 18.2 kB
  • Tags: Python 3
  • Uploaded using Trusted Publishing? Yes
  • Uploaded via: twine/6.1.0 CPython/3.13.12

File hashes

Hashes for oj_io_chains-0.0.4-py3-none-any.whl
Algorithm Hash digest
SHA256 8ab99813af6f34eee3086f1d24615bcab5fb8d59b746cc58d18ed75a5c1c45f7
MD5 dcba8eafa6df565708799e1ea21ea16f
BLAKE2b-256 825e0f7d2e27c396e07c773c03ea47a229674350809f5f73ea752221ca3edae1

See more details on using hashes here.

Provenance

The following attestation bundles were made for oj_io_chains-0.0.4-py3-none-any.whl:

Publisher: publish.yml on ownjoo-org/io-chains

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

Supported by

AWS Cloud computing and Security Sponsor Datadog Monitoring Depot Continuous Integration Fastly CDN Google Download Analytics Pingdom Monitoring Sentry Error logging StatusPage Status page