Software Engineer's Blog

Using Kafka in FastAPI — aiokafka Producer Implementation Guide

Using Kafka in FastAPI — aiokafka Producer Implementation Guide

This post is a follow-up to Kafka Infrastructure Setup & Topic Design.
Please refer to that article first for Docker Compose configuration and topic design principles.

In this post, we’ll go through how to implement a Kafka Producer in a FastAPI service.
We’ll use aiokafka directly, manage environment variables with pydantic-settings, and follow an async-first approach.

1. Dependencies (pyproject.toml)

[tool.poetry.dependencies]
aiokafka = "^0.11"
pydantic-settings = "^2.0"

2. Environment Variables

# .env (local development)
KAFKA_BOOTSTRAP_SERVERS=<KAFKA_HOST_IP>:9092
KAFKA_SECURITY_PROTOCOL=PLAINTEXT

# production (SASL_SSL)
KAFKA_BOOTSTRAP_SERVERS=kafka-broker:9092
KAFKA_SECURITY_PROTOCOL=SASL_SSL
KAFKA_SASL_MECHANISM=PLAIN
KAFKA_SASL_USERNAME=my-service
KAFKA_SASL_PASSWORD=secret

3. Define KafkaSettings

Using pydantic-settings, environment variables with the KAFKA_ prefix are automatically loaded.

# src/my_service/core/config.py
from pydantic_settings import BaseSettings, SettingsConfigDict
from pydantic import Field
from typing import Optional


class KafkaSettings(BaseSettings):
    model_config = SettingsConfigDict(env_prefix="KAFKA_", env_file=".env")

    bootstrap_servers: str = Field(..., description="Bootstrap server address")
    security_protocol: str = Field("PLAINTEXT")
    sasl_mechanism: Optional[str] = Field(None)
    sasl_username: Optional[str] = Field(None)
    sasl_password: Optional[str] = Field(None)

4. Define Event Schema

Define a BaseEvent with common fields and extend it for domain-specific events.

# src/my_service/schemas/events.py
import uuid
from datetime import datetime, timezone
from enum import Enum
from typing import Optional
from pydantic import BaseModel, Field


class BaseEvent(BaseModel):
    """Base schema for all Kafka events."""
    event_id: str = Field(default_factory=lambda: str(uuid.uuid4()))
    event_version: str = "1.0"
    occurred_at: datetime = Field(
        default_factory=lambda: datetime.now(timezone.utc)
    )
    correlation_id: Optional[str] = None
    source_service: str


class MyEventType(str, Enum):
    ITEM_CREATED = "ITEM_CREATED"
    ITEM_DELETED = "ITEM_DELETED"


class MyItemEvent(BaseEvent):
    """Payload for item create/delete events."""
    source_service: str = "my-service"
    event_type: MyEventType

    # Partition key (user identifier)
    user_id: str

    # Domain data (minimize sensitive data: include only ID and reference fields)
    item_id: str
    item_type: str

BaseEvent Common Fields

FieldTypeDescription
event_idstrUUID v4 (auto-generated)
event_versionstr"1.0" (default)
occurred_atdatetimeUTC timestamp (auto-generated)
correlation_idstrNone
source_servicestrName of the publishing service

5. KafkaProducer Wrapper

Wrap AIOKafkaProducer to handle lifecycle management and serialization.

# src/my_service/services/kafka_producer.py
import logging
from aiokafka import AIOKafkaProducer
from pydantic import BaseModel
from ..core.config import KafkaSettings

logger = logging.getLogger(__name__)


class KafkaProducer:
    def __init__(self, settings: KafkaSettings) -> None:
        kwargs: dict = {
            "bootstrap_servers": settings.bootstrap_servers,
            "security_protocol": settings.security_protocol,
        }
        if settings.security_protocol == "SASL_SSL":
            kwargs["sasl_mechanism"] = settings.sasl_mechanism
            kwargs["sasl_plain_username"] = settings.sasl_username
            kwargs["sasl_plain_password"] = settings.sasl_password

        self._producer = AIOKafkaProducer(**kwargs)

    async def start(self) -> None:
        await self._producer.start()
        logger.info("Kafka producer started")

    async def stop(self) -> None:
        await self._producer.stop()
        logger.info("Kafka producer stopped")

    async def publish(self, topic: str, event: BaseModel, key: str) -> None:
        payload = event.model_dump_json().encode("utf-8")
        key_bytes = key.encode("utf-8")
        await self._producer.send_and_wait(topic, value=payload, key=key_bytes)

6. Domain Publisher Implementation

Topic routing and error handling are handled in a domain-specific publisher class.

# src/my_service/services/kafka_publisher.py
import logging
from .kafka_producer import KafkaProducer
from ..schemas.events import MyEventType, MyItemEvent

logger = logging.getLogger(__name__)

TOPIC_MAP = {
    MyEventType.ITEM_CREATED: "myapp.my.item.created",
    MyEventType.ITEM_DELETED: "myapp.my.item.deleted",
}


class KafkaMyPublisher:
    def __init__(self, producer: KafkaProducer) -> None:
        self._producer = producer

    async def publish_item_event(self, event: MyItemEvent) -> None:
        topic = TOPIC_MAP.get(event.event_type)
        if not topic:
            raise ValueError(f"Unknown event type: {event.event_type}")

        await self._producer.publish(topic=topic, event=event, key=event.user_id)
        logger.info(
            "Kafka event published",
            extra={"topic": topic, "item_id": event.item_id, "event_type": event.event_type.value},
        )

7. Start/Stop Producer in App Lifespan

The Kafka producer should connect on startup and flush/close on shutdown.
Manage it using FastAPI’s lifespan context manager.

# src/my_service/core/app.py
from contextlib import asynccontextmanager
from fastapi import FastAPI
from .config import KafkaSettings
from .dependencies import init_kafka_publisher, shutdown_kafka_publisher


@asynccontextmanager
async def lifespan(app: FastAPI):
    # Startup: establish Kafka connection
    await init_kafka_publisher(KafkaSettings())
    yield
    # Shutdown: flush pending messages and close connection
    await shutdown_kafka_publisher()


app = FastAPI(lifespan=lifespan)

8. Manage Publisher as a Singleton

Since the producer maintains a heavy network connection, manage it as a module-level singleton.

# src/my_service/core/dependencies.py
from ..services.kafka_producer import KafkaProducer
from ..services.kafka_publisher import KafkaMyPublisher
from .config import KafkaSettings

_producer: KafkaProducer | None = None
_publisher: KafkaMyPublisher | None = None


async def init_kafka_publisher(settings: KafkaSettings) -> None:
    global _producer, _publisher
    _producer = KafkaProducer(settings)
    await _producer.start()
    _publisher = KafkaMyPublisher(_producer)


async def shutdown_kafka_publisher() -> None:
    global _producer, _publisher
    if _producer:
        await _producer.stop()
    _producer = None
    _publisher = None


def get_kafka_publisher() -> KafkaMyPublisher:
    if _publisher is None:
        raise RuntimeError("Kafka publisher not initialized")
    return _publisher

9. Publish Events in Service Layer

Use a fire-and-forget pattern when publishing events.
Event publishing should not interrupt the core business logic.

# src/my_service/services/item_service.py
import logging
from .core.dependencies import get_kafka_publisher
from .schemas.events import MyEventType, MyItemEvent

logger = logging.getLogger(__name__)


async def create_item(session, user_id: str, data: ItemCreate) -> ItemResponse:
    # 1. Save to DB
    item = await repository.create(session, data)
    await session.commit()

    # 2. Publish Kafka event (fire-and-forget)
    event = MyItemEvent(
        event_type=MyEventType.ITEM_CREATED,
        user_id=user_id,
        item_id=str(item.id),
        item_type=item.item_type,
    )
    try:
        await get_kafka_publisher().publish_item_event(event)
    except Exception:
        # Kafka failure should not affect business logic
        logger.error("Kafka publish failed, continuing", extra={"item_id": str(item.id)})

    return ItemResponse.model_validate(item)

10. Checklist for Adding Kafka to a New Service

  • [ ] Add aiokafka and pydantic-settings dependencies in pyproject.toml
  • [ ] Define KafkaSettings (auto-bind KAFKA_ environment variables)
  • [ ] Configure environment variables (.env)
  • [ ] Define domain event schema extending BaseEvent
  • [ ] Implement KafkaProducer wrapper
  • [ ] Implement KafkaMyPublisher (topic routing + error handling)
  • [ ] Add singleton management in dependencies.py
  • [ ] Initialize/shutdown in app.py lifespan
  • [ ] Publish events in service layer (fire-and-forget pattern)