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.2.tar.gz
(8.5 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.2.tar.gz.
File metadata
- Download URL: fastkafka2-0.1.2.tar.gz
- Upload date:
- Size: 8.5 kB
- Tags: Source
- Uploaded using Trusted Publishing? No
- Uploaded via: twine/6.1.0 CPython/3.13.5
File hashes
| Algorithm | Hash digest | |
|---|---|---|
| SHA256 |
729466abc4222133519fc7964d861e4d8eeab9e451c370ccf8cae7d13591b902
|
|
| MD5 |
72d1c387a214305c34f24f81b2151c1a
|
|
| BLAKE2b-256 |
6b8d0e3d67dcab82bd38bd2e2aa46ed4c02ac827b6438d733b43ee864d152f23
|
File details
Details for the file fastkafka2-0.1.2-py3-none-any.whl.
File metadata
- Download URL: fastkafka2-0.1.2-py3-none-any.whl
- Upload date:
- Size: 11.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 |
f2110018c60f0bb445db981674b132823166d89b02f9fa0a6f87a38b5d6cda16
|
|
| MD5 |
5cec4b50d919228b95de99b8768c9054
|
|
| BLAKE2b-256 |
0a2f1f86ab0387f71908124c52711b6a5c360f8ac46bf74109dd64e1ead27e74
|