Using Kafka in FastAPI — aiokafka Producer Implementation Guide
-
Jason Yang - 26 Mar, 2026
- Views —
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
| Field | Type | Description |
|---|---|---|
event_id | str | UUID v4 (auto-generated) |
event_version | str | "1.0" (default) |
occurred_at | datetime | UTC timestamp (auto-generated) |
correlation_id | str | None |
source_service | str | Name 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
aiokafkaandpydantic-settingsdependencies inpyproject.toml - [ ] Define
KafkaSettings(auto-bindKAFKA_environment variables) - [ ] Configure environment variables (
.env) - [ ] Define domain event schema extending
BaseEvent - [ ] Implement
KafkaProducerwrapper - [ ] Implement
KafkaMyPublisher(topic routing + error handling) - [ ] Add singleton management in
dependencies.py - [ ] Initialize/shutdown in
app.pylifespan - [ ] Publish events in service layer (fire-and-forget pattern)
Related Posts
- Kafka Core Concepts — Broker, Topic, Partition, Consumer Group, ACK, DLT
- Kafka Infrastructure Setup & Topic Design — Docker Compose, topic naming, common rules
- Using Kafka in Spring Boot — Spring Kafka + KafkaTemplate + @KafkaListener