Skip to main content

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


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)

Uploaded Source

Built Distribution

If you're not sure about the file name format, learn more about wheel file names.

fastkafka2-0.1.1-py3-none-any.whl (11.1 kB view details)

Uploaded Python 3

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

Hashes for fastkafka2-0.1.1.tar.gz
Algorithm Hash digest
SHA256 a3b14923203d1c3a9a13ebe655d267b9abdfb0c1aabcf08bfb0788b75212ec38
MD5 e5a6182272a4a717806f9cfcd9383ff4
BLAKE2b-256 d7fb857d9e0c2ace943c811b64634f4d4fbd587c2d6cd0eb001b2d25138d3663

See more details on using hashes here.

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

Hashes for fastkafka2-0.1.1-py3-none-any.whl
Algorithm Hash digest
SHA256 c954c1bf6e9354eac9c0c19511e94f129c845dd738d339cb5d6cf02335787816
MD5 34fc22289252217003eb83aeae3bf9a8
BLAKE2b-256 8cd5dcb8e50b0a589e6943c294d7959ff7b71f0d5a715a817cafa5a0a0131c40

See more details on using hashes here.

Supported by

AWS Cloud computing and Security Sponsor Datadog Monitoring Depot Continuous Integration Fastly CDN Google Download Analytics Pingdom Monitoring Sentry Error logging StatusPage Status page