Next-generation FastAPI-like DX for Kafka (version 2)
Project description
fastkafka2
Next-generation FastAPI-like DX for Kafka (version 2).
Installation
pip install fastkafka2
Use example
Project arch
├── api/
│ ├── kafka/
│ │ ├── handlers/
│ │ │ ├── example/
│ │ │ │ ├─ schemas.py
│ │ │ │ └─ handler.py
│ │ │ └── base_handler.py
│ │ └── lifespan.py
│
├── main.py
Handler
# api\kafka\handlers\example\handler.py
import logging
from fastkafka2 import KafkaHandler
from fastkafka2 import KafkaMessage
from fastkafka2 import KafkaProducer
handler = KafkaHandler()
kafka_producer = KafkaProducer(bootstrap_servers="127.0.0.1:9092")
@handler("example")
async def example_handler(message: KafkaMessage):
t = int(message.headers.get("try")) + 1
logging.info(f"Пришло: {message}")
await kafka_producer.send_message(
topic="example-2", data={"msg": "wddwd"}, headers={"try": f"{t}"}, key=None
)
logging.info(f"Отправил: {f'{t}'}")
Typed validation and IDE hints
There are two ways to get strong typing and validation like in FastAPI:
- Single-source via function annotation (recommended for IDE hints)
from pydantic import BaseModel
from fastkafka2 import KafkaHandler
from fastkafka2 import KafkaMessage
class Data(BaseModel):
id: int
amount: float
class Headers(BaseModel):
type: str
source: str
handler = KafkaHandler(prefix="orders")
@handler("created") # models inferred from annotation
async def on_created(msg: KafkaMessage[Data, Headers]):
# msg.data and msg.headers are fully typed
...
Grouping of handlers
# api\kafka\handlers\base_handler.py
from api.kafka.handlers.example.handler import handler as example_handler
from fastkafka2 import KafkaHandler
base_handler = KafkaHandler()
base_handler.include_handler(example_handler)
Lifespan fastkafka app
# api/kafka/lifespan.py
import logging
from contextlib import asynccontextmanager
from fastkafka2 import KafkaApp
from api.kafka.handlers.base_handler import base_handler
from api.kafka.handlers.example.handler import kafka_producer
@asynccontextmanager
async def lifespan(app: KafkaApp):
logging.info("Lifespan: запуск")
try:
await kafka_producer.start()
yield
logging.info("Lifespan: выполнен")
finally:
await kafka_producer.stop()
logging.info("Lifespan: остановка")
app = KafkaApp(
title="Kafka Gateway",
description="Kafka-based microservice",
bootstrap_servers="127.0.0.1:9092",
group_id="test_service",
lifespan=lifespan,
)
app.include_handler(base_handler)
Entry point main app
# main.py
import asyncio
from logging_config import setup_logging
from api.kafka.lifespan import app
if __name__ == "__main__":
setup_logging()
asyncio.run(app.run())
Project details
Release history Release notifications | RSS feed
Download files
Download the file for your platform. If you're not sure which to choose, learn more about installing packages.
Source Distribution
fastkafka2-0.1.5.tar.gz
(9.4 kB
view details)
Built Distribution
Filter files by name, interpreter, ABI, and platform.
If you're not sure about the file name format, learn more about wheel file names.
Copy a direct link to the current filters
File details
Details for the file fastkafka2-0.1.5.tar.gz.
File metadata
- Download URL: fastkafka2-0.1.5.tar.gz
- Upload date:
- Size: 9.4 kB
- Tags: Source
- Uploaded using Trusted Publishing? No
- Uploaded via: twine/6.2.0 CPython/3.14.0
File hashes
| Algorithm | Hash digest | |
|---|---|---|
| SHA256 |
6452ba8280d6400eee30e5816014a439538939298ea6f5a98d6d40e4903479aa
|
|
| MD5 |
adef9529e8ce4f87c3b0ae30cccba020
|
|
| BLAKE2b-256 |
4a17573032c84d3f91ef1b5488d0b447b48528bbda337a3aaba880a2ee4abfa7
|
File details
Details for the file fastkafka2-0.1.5-py3-none-any.whl.
File metadata
- Download URL: fastkafka2-0.1.5-py3-none-any.whl
- Upload date:
- Size: 12.4 kB
- Tags: Python 3
- Uploaded using Trusted Publishing? No
- Uploaded via: twine/6.2.0 CPython/3.14.0
File hashes
| Algorithm | Hash digest | |
|---|---|---|
| SHA256 |
5fd2a071f91fd4f61f68ac892ca9d4c2eb2b25779216f76f84f6f3ff7ce4ecd8
|
|
| MD5 |
db9ee6bd0c6081b813930f71ddc9021e
|
|
| BLAKE2b-256 |
fbacb6b4a148560f4fabdec64884878f49e1f1aee1bbfb59bdafc3258f83442c
|