Skip to main content

threaded-buffered-pipeline CircleCI Test Coverage

Parallelise pipelines of Python iterables

Installation

pip install threaded-buffered-pipeline

Usage / What problem does this solve?

If you have a chain of generators, only one runs at any given time. For example, the below runs in (just over) 30 seconds.

import time

def gen_1():
    for value in range(0, 10):
        time.sleep(1)  # Could be a slow HTTP request
        yield value

def gen_2(it):
    for value in it:
        time.sleep(1)  # Could be a slow HTTP request
        yield value * 2

def gen_3(it):
    for value in it:
        time.sleep(1)  # Could be a slow HTTP request
        yield value + 3

def main():
    it_1 = gen_1()
    it_2 = gen_2(it_1)
    it_3 = gen_3(it_2)

    for val in it_3:
        print(val)

main()

The buffered_pipeline function allows you to make to a small change, passing each generator through its return value, to parallelise the generators to reduce this to (just over) 12 seconds.

import time
from threaded_buffered_pipeline import buffered_pipeline

def gen_1():
    for value in range(0, 10):
        time.sleep(1)  # Could be a slow HTTP request
        yield value

def gen_2(it):
    for value in it:
        time.sleep(1)  # Could be a slow HTTP request
        yield value * 2

def gen_3(it):
    for value in it:
        time.sleep(1)  # Could be a slow HTTP request
        yield value + 3

def main():
    buffer_iterable = buffered_pipeline()
    it_1 = buffer_iterable(gen_1())
    it_2 = buffer_iterable(gen_2(it_1))
    it_3 = buffer_iterable(gen_3(it_2))

    for val in it_3:
        print(val)

main()

The buffered_pipeline ensures internal threads are stopped on any exception [the next time each thread attempts to pull from the iterator].

Buffer size

The default buffer size is 1. This is suitable if each iteration takes approximately the same amount of time. If this is not the case, you may wish to change it using the buffer_size parameter of buffer_iterable.

it = buffer_iterable(gen(), buffer_size=2)

Features

  • One thread is created for each buffer_iterable, in which the iterable is iterated over, with its values stored in an internal buffer.

  • All the threads of the pipeline are stopped if any of the generators raise an exception.

  • If a generator raises an exception, the exception is propagated to calling code.

  • The buffer size of each step in the pipeline is configurable.

  • The "chaining" is not abstracted away. You still have full control over the arguments passed to each step, and you don't need to buffer each iterable in the pipeline if you don't want to: just don't pass those through buffer_iterable.

Asyncio

A version for async iterables is available at https://github.com/michalc/asyncio-buffered-pipeline

Metadata

Release files for threaded-buffered-pipeline 0.0.8

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

Source distribution (sdist)

Source distribution for threaded-buffered-pipeline 0.0.8
File Size Uploaded
threaded-buffered-pipeline-0.0.8.tar.gz 3.7 kB Details

Built distribution (wheel)

Table of built distributions (wheels) for threaded-buffered-pipeline 0.0.8
File Interpreter ABI Platform
threaded_buffered_pipeline-0.0.8-py3-none-any.whl Python 3 none any Details

Total release size: 8.3 kB

Release files / threaded-buffered-pipeline-0.0.8.tar.gz

Download URL threaded-buffered-pipeline-0.0.8.tar.gz
Size 3.7 kB
Tags Source
SHA-256 checksum
How to use checksums
4b620cb56cb166ee83b0284422e37ddf764a727ffddbe04d3c5f43e6b79620b0
BLAKE2b-256 checksum
How to use checksums
59844c87a9da69eed0ff98eabe1a36cd91c52e8a5240fcc143bb152460824be6
Upload date
Uploaded using Trusted Publishing?
What is trusted publishing?
No
Uploaded via twine/3.2.0 pkginfo/1.5.0.1 requests/2.24.0 setuptools/50.3.2 requests-toolbelt/0.9.1 tqdm/4.49.0 CPython/3.8.5

Release files / threaded_buffered_pipeline-0.0.8-py3-none-any.whl

Download URL threaded_buffered_pipeline-0.0.8-py3-none-any.whl
Size 4.7 kB
Tags Python 3
SHA-256 checksum
How to use checksums
709318e5b8e52802b5c269d7166248dd10a7d4d4f239e5dd939c65e7358a3099
BLAKE2b-256 checksum
How to use checksums
0d2e99503cbb378b29d1ecfd82de68f2d37ab9c3a6e4a25367cd0b6c3cb908d7
Upload date
Uploaded using Trusted Publishing?
What is trusted publishing?
No
Uploaded via twine/3.2.0 pkginfo/1.5.0.1 requests/2.24.0 setuptools/50.3.2 requests-toolbelt/0.9.1 tqdm/4.49.0 CPython/3.8.5

Release history Release notifications | RSS feed

This release

0.0.8 This release

2 release files

0.0.7

2 release files

0.0.6

2 release files

0.0.5

2 release files

0.0.4

2 release files

0.0.3

2 release files

0.0.2

2 release files

0.0.1

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