Skip to main content

Nonblocking Stream Queue

This is a simple python package for nonblocking reading from streams.

It could be expanded for general nonblocking usage if desired.

A simple nonblocking queue reader is provided, with implementation from https://stackoverflow.com/questions/375427/a-non-blocking-read-on-a-subprocess-pipe-in-python/4896288#4896288 .

When constructed, the reader spawns a thread and begins reading everything from the stream.

Installation

$ python3 -m pip install nonblocking-stream-queue

Usage

import sys
import nonblocking_stream_queue as nonblocking

reader = nonblocking.Reader(
    sys.stdin, # object to read from
    max_size=4096, # max size of each read
    lines=False, # whether or not to break reads into lines
    max_count=None, # max queued reads
    drop_timeout=None, # time to wait for queue to drain before dropping
    drop_older=False, # which end of the queue to drop from
    pre_cb=None, # if set, data = (pre_cb(), read())
    post_cb=None, # if set, data = post_cb(data)
    drop_cb=None, # if set, call with dropped data
    verbose=False, # if set, display progress filling max_count
)

print(reader.read_one()) # outputs up to 4096 characters, or None if nothing queued
print(reader.read_many()) # outputs all 4096 character chunks queued
print(''.join(reader.read_many())) # outputs all queued text
reader.block() # waits for data to be present or end, returns number of reads queued
reader.is_pumping() # False if data has ended

Timestamping

import sys, time
import nonblocking_stream_queue as nonblocking

reader = nonblocking.Reader(
    sys.stdin.buffer,
    pre_cb=lambda: time.time(),
    post_cb=lambda time_data_tuple: (time_data_tuple[0], time_data_tuple[1], time.time())
)

while reader.block():
    start_time, data, end_time = reader.read_one()
    print(start_time, data, end_time)

Lines

import sys
import nonblocking_stream_queue as nonblocking

reader = nonblocking.Reader(
    sys.stdin,
    lines=True
)

while reader.block():
    print(reader.read_one().rstrip()) # outputs each line of text that is input

Other solutions

There are likely many other existing solutions to this common task.

Here's one:

Release files for nonblocking-stream-queue 0.1.4

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

Source distribution (sdist)

Source distribution for nonblocking-stream-queue 0.1.4
File Size Uploaded
nonblocking_stream_queue-0.1.4.tar.gz 3.6 kB Details

Built distribution (wheel)

Table of built distributions (wheels) for nonblocking-stream-queue 0.1.4
File Interpreter ABI Platform
nonblocking_stream_queue-0.1.4-py3-none-any.whl Python 3 none any Details

Total release size:8.0 kB

Release files / nonblocking_stream_queue-0.1.4.tar.gz

Download URL nonblocking_stream_queue-0.1.4.tar.gz
Size 3.6 kB
Tags Source
SHA-256 checksum
How to use checksums
fd3cf4112b3a194984a4078a2c588222a389ee1ad58727817c3d2a656aa5907f
BLAKE2b-256 checksum
How to use checksums
4fa79dbf258b9ad19706ba14ea5b5af17567518ec105e4fba02dffd1346c7030
Upload date
Uploaded using Trusted Publishing?
What is trusted publishing?
No
Uploaded via twine/4.0.0 CPython/3.10.4+

Release files / nonblocking_stream_queue-0.1.4-py3-none-any.whl

Download URL nonblocking_stream_queue-0.1.4-py3-none-any.whl
Size 4.4 kB
Tags Python 3
SHA-256 checksum
How to use checksums
7f4c10303648f33272a31292e2118e86535bf427fea63a1db463ba30c835655d
BLAKE2b-256 checksum
How to use checksums
34200c66086b5492e15bf0d10d8f162b678f2f7530b63ea46c9af76ab21ecbc9
Upload date
Uploaded using Trusted Publishing?
What is trusted publishing?
No
Uploaded via twine/4.0.0 CPython/3.10.4+

Release history Release notifications | RSS feed

This release

0.1.4 This release

2 release files

0.1.3

2 release files

0.1.1

2 release files

0.1.0

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