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)
| File | Size | Uploaded | |
|---|---|---|---|
| bigpipe-0.1.0.tar.gz | 12.7 kB | Details |
Built distribution (wheel)
| File | Interpreter | ABI | Platform | Reset |
|---|---|---|---|---|
| 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
|