Skip to main content

StreamMachine

Async stream processing on Redis Streams for Python. Register consumers and periodic producers with decorators, run them on one event loop, and scale out by starting more processes in the same consumer group.

from streammachine import App, Message

app = App(name="demo")

@app.timer(1)
async def producer():
    await app.send("greetings", {"message": "hello"})

@app.agent("greetings", group="greeters")
async def consumer(record: Message):
    print("received:", record.message)

if __name__ == "__main__":
    app.start()

Features

  • Agents and timers: @app.agent(stream, group=...) consumes a stream through a Redis consumer group, @app.timer(seconds) runs a coroutine periodically.
  • Resilient consumers: handler errors are logged and skipped; Redis outages reconnect with bounded exponential backoff under a stable consumer name so pending entries are never orphaned.
  • Bounded streams: App(stream_maxlen=...) or STREAMMACHINE_STREAM_MAXLEN trims every produced stream with approximate MAXLEN.
  • Shared state: app.storage is a multiprocessing.Manager backed key/value store with per-key async locks.
  • DataFrames: streams_to_dataframe() and TimeSeriesBuffer turn raw XREAD output into pandas frames with automatic pruning.
  • OHLC aggregation: create_ohlc_aggregator() builds candles from tick streams, with an optional Cython build for higher throughput.
  • Optional extras: a FastAPI monitoring dashboard, an MCP server exposing streams and storage as tools, and pickle-based object storage.

Installation

pip install streammachine

Extras:

Extra Installs Purpose
streammachine[dashboard] fastapi, uvicorn Web dashboard on http://localhost:8000
streammachine[mcp] mcp streammachine-mcp server for LLM clients
streammachine[objstorage] redis RedisObjectStorage (pickle in Redis)
streammachine[cython] cython Build accelerators with STREAMMACHINE_BUILD_CYTHON=1
streammachine[all] dashboard, mcp, objstorage Everything above except Cython

Requires Python 3.10+ and a Redis server (6.2 or newer recommended). uvloop is used automatically on Linux and macOS.

Configuration

Connection settings are read from the environment:

Variable Default Meaning
REDIS_URL redis://localhost:6379 Full connection URL
REDIS_HOST / REDIS_PORT / REDIS_DB localhost / 6379 / 0 Used when no URL is given
REDIS_MAX_CONNECTIONS 10 Pool size per connection
STREAMMACHINE_DEFAULT_GROUP eventengine Consumer group when group= is omitted
STREAMMACHINE_STREAM_MAXLEN unset Approximate max length for produced streams

How it works

  1. Decorators attach metadata to your handlers (via venusian); nothing runs at import time.
  2. app.start() scans the calling module, creates one consumer task per agent (times concurrency), one task per timer, and starts the shared storage manager.
  3. Each consumer joins its consumer group with XREADGROUP, wraps every entry in a Message (topic, stream id, decoded fields, send/receive timestamps) and awaits your handler.
  4. SIGINT/SIGTERM trigger a graceful shutdown: timers stop, tasks are cancelled with a timeout, Redis connections and the storage manager are closed.

Run several copies of the same script to scale horizontally; Redis distributes entries across consumers in a group.

Documentation

Development

git clone https://github.com/trbck/streammachine.git
cd streammachine
python -m venv .venv && source .venv/bin/activate
pip install -e ".[dev,all]"
pytest                                  # unit tests (no Redis needed)
RUN_INTEGRATION_TESTS=1 pytest          # integration tests against a local Redis
ruff check src tests

Build and check the distribution:

python -m build
twine check dist/*

License

Apache License 2.0. See LICENSE.

Metadata

Release files for streammachine 0.1.0

For a detailed explanation of source distributions (sdists) and built distributions (wheels), please see the package formats documentation.

Source distribution (sdist)

Source distribution for streammachine 0.1.0
File Size Uploaded
streammachine-0.1.0.tar.gz 105.4 kB Details

Built distribution (wheel)

Table of built distributions (wheels) for streammachine 0.1.0
File Interpreter ABI Platform
streammachine-0.1.0-py3-none-any.whl Python 3 none any Details

Total release size: 182.5 kB

Release files / streammachine-0.1.0.tar.gz

Download URL streammachine-0.1.0.tar.gz
Size 105.4 kB
Tags Source
SHA-256 checksum
How to use checksums
3314b7945b2c1343d4162b582eeddb419194273aadeeb5cfa38e117b20c56624
BLAKE2b-256 checksum
How to use checksums
bd8a5165a2dce5ea5fb23f28006ac761d96233f39ea9ce1cb47fb0889803f968
Upload date
Uploaded using Trusted Publishing?
What is trusted publishing?
Yes
Uploaded via twine/7.0.0 CPython/3.13.14

Provenance

Provenance describes where a file came from. On PyPI, provenance is shared via attestations, which provide a verifiable record of the build or publishing details. View details, limitations and caveats.

PyPI Publish Attestation

PyPI verified that this artifact, at this checksum, originated from the publisher listed below.

Signed by GitHub Actions, verified by PyPI on Sep 26, 2026.

Transparency log

Release files / streammachine-0.1.0-py3-none-any.whl

Download URL streammachine-0.1.0-py3-none-any.whl
Size 77.1 kB
Tags Python 3
SHA-256 checksum
How to use checksums
c4fee9d224e6e1683d1b9bd8adc87e911ce51366531c72a357cdc499c07de999
BLAKE2b-256 checksum
How to use checksums
a93086a847217324592fa85fb300d2ff884b993ae862600bc2071652d5df1d5f
Upload date
Uploaded using Trusted Publishing?
What is trusted publishing?
Yes
Uploaded via twine/7.0.0 CPython/3.13.14

Provenance

Provenance describes where a file came from. On PyPI, provenance is shared via attestations, which provide a verifiable record of the build or publishing details. View details, limitations and caveats.

PyPI Publish Attestation

PyPI verified that this artifact, at this checksum, originated from the publisher listed below.

Signed by GitHub Actions, verified by PyPI on Sep 26, 2026.

Transparency log

Release history Release notifications | RSS feed

This release

0.1.0 This release

2 release 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