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
Schemas
# api\kafka\handlers\example\schemas.py
from pydantic import BaseModel
class ExampleSchema(BaseModel):
msg: str
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 Order(BaseModel):
id: int
amount: float
class Hdrs(BaseModel):
type: str
source: str
handler = KafkaHandler(prefix="orders")
@handler("created") # models inferred from annotation
async def on_created(msg: KafkaMessage[Order, Hdrs]):
# msg.data and msg.headers are fully typed
...
- Models in decorator, generic message in function (runtime validation only)
@handler("created", data_model=Order, headers_model=Hdrs)
async def on_created(msg: KafkaMessage):
# Works at runtime; for IDE hints you can optionally:
# from typing import cast
# msg = cast(KafkaMessage[Order, Hdrs], msg)
...
You can also split parameters:
@handler("updated")
async def on_updated(data: Order, headers: Hdrs):
...
Header filtering
Filter messages by headers before deserialization using equality or a predicate:
# Equality filter
@handler("created", data_model=Order, headers_model=Hdrs, headers_filter={"type": "created"})
async def on_created(msg: KafkaMessage[Order, Hdrs]):
...
# Predicate filter
def high_value(h: dict[str, str]) -> bool:
return h.get("type") == "created" and h.get("priority") == "high"
@handler("created", data_model=Order, headers_model=Hdrs, headers_filter=high_value)
async def on_high_value(msg: KafkaMessage[Order, Hdrs]):
...
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",
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.4.tar.gz
(10.0 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.4.tar.gz.
File metadata
- Download URL: fastkafka2-0.1.4.tar.gz
- Upload date:
- Size: 10.0 kB
- Tags: Source
- Uploaded using Trusted Publishing? No
- Uploaded via: twine/6.1.0 CPython/3.13.5
File hashes
| Algorithm | Hash digest | |
|---|---|---|
| SHA256 |
667eb26ab901a0a584e72df0ec6429d7682d240c702333fcac26a4ad4f953944
|
|
| MD5 |
afeb7cb44ec65f7fda1c255f86e09750
|
|
| BLAKE2b-256 |
b6b2e76eb5f04102be8f19809bc6acc3bde73c5f008e781444b6ae982fb9bdff
|
File details
Details for the file fastkafka2-0.1.4-py3-none-any.whl.
File metadata
- Download URL: fastkafka2-0.1.4-py3-none-any.whl
- Upload date:
- Size: 12.7 kB
- Tags: Python 3
- Uploaded using Trusted Publishing? No
- Uploaded via: twine/6.1.0 CPython/3.13.5
File hashes
| Algorithm | Hash digest | |
|---|---|---|
| SHA256 |
b590f84f0b3d75f11ebb048de71536ac2a512586d652e4527e52e2440d7a7657
|
|
| MD5 |
7db0d6e4f3014ceca090300d9c304189
|
|
| BLAKE2b-256 |
27eeea22bf9750931c54e5e71e0c0fad4db1e81e8b602ef329036dad27936ff9
|