Skip to main content

bigpipe — Python SDK for BigPipe

Async Python client for BigPipe, the high performance realtime stream processing platform.

Dibuat oleh Gravicode Studios dipimpin oleh Kang Fadhil.

pip install bigpipe

Produce and consume

import asyncio
from bigpipe import Producer, Consumer

async def main():
    async with Producer(bootstrap="localhost:9092") as p:
        meta = await p.send("orders", {"id": 1, "amount": 150000}, key="order-1", headers={"source": "checkout"})
        print(f"written to {meta.partition}@{meta.offset}")

    async with Consumer(bootstrap="localhost:9092", group="analytics", topics=["orders"],
                        auto_offset_reset="earliest") as c:
        async for msg in c:
            print(msg.key_str, msg.json())
            await c.commit(msg)

asyncio.run(main())

Producer batches everything sent within linger_ms into one request, so await asyncio.gather(*(p.send(...) for ...)) is fast. Values may be bytes, str or anything JSON-serializable.

Share groups (work queues with a dead-letter topic)

from bigpipe import ShareConsumer

async with ShareConsumer(group="email-workers", topics=["email-jobs"],
                         max_attempts=5, dlq_topic="email-jobs.dlq") as worker:
    async for job in worker:
        try:
            await send_email(job.json())
            job.accept()
        except TemporaryError:
            job.release()   # redelivered to any worker
        except Exception:
            job.reject()    # goes to the dead-letter topic

Streaming with a server-side filter

from bigpipe import stream

async for r in stream("payments", from_="latest", filter='header("region") == "ID" && this.amount > 1000000'):
    print(r.json())

Admin

from bigpipe import Admin

async with Admin("http://localhost:9644") as admin:
    await admin.create_topic("clicks", partitions=12, mode="diskless")
    await admin.migrate("orders", "tiered")          # online, offsets unchanged
    print(await admin.group("analytics"))           # lag per partition

The SDK talks to the BigPipe HTTP gateway (default port 8082) and admin API (9644). bootstrap="host:9092" maps to http://host:8082; pass url= to override. Environment variables: BIGPIPE_HTTP, BIGPIPE_ADMIN, BIGPIPE_API_KEY.

License: Apache-2.0.

Metadata

Release files for bigpipe 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 bigpipe 0.1.0
File Size Uploaded
bigpipe-0.1.0.tar.gz 12.7 kB Details

Built distribution (wheel)

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

Total release size: 25.2 kB

Release files / bigpipe-0.1.0.tar.gz

Download URL bigpipe-0.1.0.tar.gz
Size 12.7 kB
Tags Source
SHA-256 checksum
How to use checksums
a6d6a60961a67c6fc09293fcf5b208837c851adeb189954c426b75e10122cf2c
BLAKE2b-256 checksum
How to use checksums
8b0f065b8443be142eaf3fb7a146f97638acf333ee67d4695322976ec1d2544b
Upload date
Uploaded using Trusted Publishing?
What is trusted publishing?
No
Uploaded via twine/7.0.0 CPython/3.12.10

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

Download URL bigpipe-0.1.0-py3-none-any.whl
Size 12.5 kB
Tags Python 3
SHA-256 checksum
How to use checksums
304295c84f45e6867faf06c253013a22ca747eb8f19e48a053f304a240eb9e84
BLAKE2b-256 checksum
How to use checksums
c09a2427238cd27ebfe49aa4857bdee62f754b67346aae4e1dc3bb6ee9c6aea8
Upload date
Uploaded using Trusted Publishing?
What is trusted publishing?
No
Uploaded via twine/7.0.0 CPython/3.12.10

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