Skip to main content

StreamPipe

StreamPipe is a simple Python library (a single file!) for running multi-worker pipelines on infinite data streams. The best part about StreamPipe is that the order of data is preserved throughout the pipeline. You can chain order-independent workers (think of image processing functions) and order-dependent workers (think of an object video tracker) on the same StreamPipe!

Installation

pip3 install streampipe

Install from source:

git clone https://github.com/highvight/streampipe
cd streampipe
pip3 install .
pip3 install .[testing]  # for testing and benchmarks

Usage

Here is a starter. Let's check if you can speed up a reshaping and standardization pipeline for images.

import timeit
import numpy as np
from streampipe import StreamPipe

IMG = np.random.rand(1000, 1000, 3) * 255

def _reshape(img):
    return img[::2, ::2]

def _standardize(img):
    mean = np.mean(img, axis=(0, 1), keepdims=True)
    std = np.std(img, axis=(0, 1), keepdims=True)
    return (img - mean) / std

pipe = StreamPipe()
pipe.add(_reshape, n_workers=4, maxsize=4)
pipe.add(_standardize, n_workers=4, maxsize=4)

def measure_pipe():
    # Starts the workers
    with pipe:
        for _ in range(100):
            pipe.put(IMG)
    # All workers stopped, all items done

def measure_normal():
    for _ in range(100):
        x = _reshape(IMG)
        x = _standardize(x)

print(timeit.timeit(measure_pipe, number=1))  # 1.214164252000046
print(timeit.timeit(measure_normal, number=1))  # 3.083142001007218

License

MIT

Release files for streampipe 0.1

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

Source distribution (sdist)

Source distribution for streampipe 0.1
File Size Uploaded
streampipe-0.1.tar.gz 6.9 kB Details

Built distribution (wheel)

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

Total release size: 12.9 kB

Release files / streampipe-0.1.tar.gz

Download URL streampipe-0.1.tar.gz
Size 6.9 kB
Tags Source
SHA-256 checksum
How to use checksums
be493d7f50ac33874d7b4e42ced152f7198773f39f2038d5a3b29eebfbfa15fb
BLAKE2b-256 checksum
How to use checksums
a50ab32cc52d5ef7482ea68d471c5eb7ddbdbcfb34b5d662da652954593c8c53
Upload date
Uploaded using Trusted Publishing?
What is trusted publishing?
No
Uploaded via twine/4.0.2 CPython/3.9.18

Release files / streampipe-0.1-py3-none-any.whl

Download URL streampipe-0.1-py3-none-any.whl
Size 6.0 kB
Tags Python 3
SHA-256 checksum
How to use checksums
129dd135034a59797e1ba0d8150d0b2b1e1776f78a00e1eec9f38b83012c1ec0
BLAKE2b-256 checksum
How to use checksums
902eac6e62a9be125e400d4c41ca22b67bd5c315c2930dd9d8d11564275c01b7
Upload date
Uploaded using Trusted Publishing?
What is trusted publishing?
No
Uploaded via twine/4.0.2 CPython/3.9.18

Release history Release notifications | RSS feed

This release

0.1 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