Kstreams
kstreams is a library/micro framework to use with kafka. It has simple kafka streams implementation that gives certain guarantees, see below.
Documentation: https://kpn.github.io/kstreams/
Installation
pip install kstreams
You will need a worker, we recommend aiorun
pip install aiorun
Usage
import aiorun
from kstreams import create_engine, ConsumerRecord
stream_engine = create_engine(title="my-stream-engine")
@stream_engine.stream("local--kstream")
async def consume(cr: ConsumerRecord):
print(f"Event consumed: headers: {cr.headers}, payload: {cr.value}")
async def produce():
payload = b'{"message": "Hello world!"}'
for i in range(5):
metadata = await stream_engine.send("local--kstreams", value=payload)
print(f"Message sent: {metadata}")
async def start():
await stream_engine.start()
await produce()
async def shutdown(loop):
await stream_engine.stop()
if __name__ == "__main__":
aiorun.run(start(), stop_on_unhandled_errors=True, shutdown_callback=shutdown)
Features
- Produce events
- Consumer events with
Streams - Subscribe to topics by
pattern -
Prometheusmetrics and custom monitoring - TestClient
- Custom Serialization and Deserialization
- Easy to integrate with any
asyncframework. No tied to any library!! - Yield events from streams
- Opentelemetry Instrumentation
- Middlewares
- Hooks (on_startup, on_stop, after_startup, after_stop)
- Store (kafka streams pattern)
- Stream Join
- Windowing
Development
This repo requires the use of poetry instead of pip.
Note: If you want to have the virtualenv in the same path as the project first you should run poetry config --local virtualenvs.in-project true
To install the dependencies just execute:
poetry install
Then you can activate the virtualenv with
poetry shell
Run test:
./scripts/test
Run code formatting with ruff:
./scripts/format
Commit messages
We use conventional commits for the commit message.
The use of commitizen is recommended. Commitizen is part of the dev dependencies.
cz commit
Metadata
Release files for kstreams 0.34.3
For a detailed explanation of source distributions (sdists) and built distributions (wheels), please see the package formats documentation.
Source distribution (sdist)
| File | Size | Uploaded | |
|---|---|---|---|
| kstreams-0.34.3.tar.gz | 36.0 kB | Details |
Built distribution (wheel)
| File | Interpreter | ABI | Platform | Reset |
|---|---|---|---|---|
| kstreams-0.34.3-py3-none-any.whl | Python 3 | none | any | Details |
Total release size: 81.5 kB
Release files / kstreams-0.34.3.tar.gz
| Download URL | kstreams-0.34.3.tar.gz |
|---|---|
| Size | 36.0 kB |
| Tags | Source |
|
SHA-256 checksum How to use checksums |
49c07d9b453e2115ef0aa7bc8cb444d6ea46986710a238a5713debd0cc152e35
|
|
BLAKE2b-256 checksum How to use checksums |
7416451c43fd502bde03022d4e9441ce9ad027bf2d4f14ad0f353261af393ffd
|
| Upload date | |
|
Uploaded using Trusted Publishing? What is trusted publishing? |
No |
| Uploaded via |
poetry/2.4.3 CPython/3.14.7 Linux/6.17.0-1022-azure
|
Release files / kstreams-0.34.3-py3-none-any.whl
| Download URL | kstreams-0.34.3-py3-none-any.whl |
|---|---|
| Size | 45.6 kB |
| Tags | Python 3 |
|
SHA-256 checksum How to use checksums |
3f022604b6d3fb4606cc0aab03dcb69ab491ce62b8534a91a188c07f5f18d776
|
|
BLAKE2b-256 checksum How to use checksums |
684d39f0e88a7221cc1908a282a4f7f35d0480460dca47335f990624e673a554
|
| Upload date | |
|
Uploaded using Trusted Publishing? What is trusted publishing? |
No |
| Uploaded via |
poetry/2.4.3 CPython/3.14.7 Linux/6.17.0-1022-azure
|