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 fastkafka.handler import KafkaHandler
from fastkafka.message import KafkaMessage
from fastkafka.producer 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}'}")
Grouping of handlers
# api\kafka\handlers\base_handler.py
from api.kafka.handlers.example.handler import handler as example_handler
from fastkafka.handler 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 fastkafka.app 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.1.tar.gz
(7.9 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.1.tar.gz.
File metadata
- Download URL: fastkafka2-0.1.1.tar.gz
- Upload date:
- Size: 7.9 kB
- Tags: Source
- Uploaded using Trusted Publishing? No
- Uploaded via: twine/6.1.0 CPython/3.13.5
File hashes
| Algorithm | Hash digest | |
|---|---|---|
| SHA256 |
a3b14923203d1c3a9a13ebe655d267b9abdfb0c1aabcf08bfb0788b75212ec38
|
|
| MD5 |
e5a6182272a4a717806f9cfcd9383ff4
|
|
| BLAKE2b-256 |
d7fb857d9e0c2ace943c811b64634f4d4fbd587c2d6cd0eb001b2d25138d3663
|
File details
Details for the file fastkafka2-0.1.1-py3-none-any.whl.
File metadata
- Download URL: fastkafka2-0.1.1-py3-none-any.whl
- Upload date:
- Size: 11.1 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 |
c954c1bf6e9354eac9c0c19511e94f129c845dd738d339cb5d6cf02335787816
|
|
| MD5 |
34fc22289252217003eb83aeae3bf9a8
|
|
| BLAKE2b-256 |
8cd5dcb8e50b0a589e6943c294d7959ff7b71f0d5a715a817cafa5a0a0131c40
|