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.2.tar.gz (8.5 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.2-py3-none-any.whl (11.7 kB view details)

Uploaded Python 3

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

Hashes for fastkafka2-0.1.2.tar.gz
Algorithm Hash digest
SHA256 729466abc4222133519fc7964d861e4d8eeab9e451c370ccf8cae7d13591b902
MD5 72d1c387a214305c34f24f81b2151c1a
BLAKE2b-256 6b8d0e3d67dcab82bd38bd2e2aa46ed4c02ac827b6438d733b43ee864d152f23

See more details on using hashes here.

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

Hashes for fastkafka2-0.1.2-py3-none-any.whl
Algorithm Hash digest
SHA256 f2110018c60f0bb445db981674b132823166d89b02f9fa0a6f87a38b5d6cda16
MD5 5cec4b50d919228b95de99b8768c9054
BLAKE2b-256 0a2f1f86ab0387f71908124c52711b6a5c360f8ac46bf74109dd64e1ead27e74

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