Kafka с aiokafka: topics, partitions, consumer groups, offsets
Kafka с aiokafka: topics, partitions, consumer groups, offsets
В отличие от RabbitMQ с его queues, Kafka это distributed append-only log. Сообщения хранятся в topics разбитых на partitions, потребляются consumer groups с tracking offsets. Это другая mental model - в этом уроке разберём как работать с Kafka из Python через aiokafka.
Установка и запуск
pip install aiokafka
Kafka через docker-compose:
version: "3"
services:
zookeeper:
image: confluentinc/cp-zookeeper:7.6.0
environment:
ZOOKEEPER_CLIENT_PORT: 2181
kafka:
image: confluentinc/cp-kafka:7.6.0
depends_on: [zookeeper]
ports:
- "9092:9092"
environment:
KAFKA_BROKER_ID: 1
KAFKA_ZOOKEEPER_CONNECT: zookeeper:2181
KAFKA_ADVERTISED_LISTENERS: PLAINTEXT://localhost:9092
KAFKA_OFFSETS_TOPIC_REPLICATION_FACTOR: 1
docker-compose up -d
С Kafka 3.0+ можно без ZooKeeper (KRaft mode), но docker-compose всё ещё проще с ZK.
Концепты
- Topic - named stream (например,
user-events) - Partition - ordered subset topic. Topic = N partitions
- Offset - позиция в partition (0, 1, 2, ...)
- Producer - отправляет messages в topic
- Consumer - читает с определённого offset
- Consumer group - набор consumers разделяющих работу по partitions (подробный разбор ребалансировки - в треке message-brokers)
Каждое message → один partition (по key или round-robin). Внутри partition - strict order, между partitions - нет.
Producer
import asyncio
from aiokafka import AIOKafkaProducer
import json
async def produce():
producer = AIOKafkaProducer(bootstrap_servers="localhost:9092")
await producer.start()
try:
for i in range(10):
data = {"event": "user_login", "user_id": i}
await producer.send_and_wait(
"user-events",
value=json.dumps(data).encode(),
key=str(i).encode(), # для partition routing
)
print("Sent 10 messages")
finally:
await producer.stop()
asyncio.run(produce())
bootstrap_servers- начальный broker (один достаточно, остальные discover автоматически)valueиkey- bytes (нужно serialize)send_and_wait- sync semantic (ждёт подтверждение)send- fire-and-forget, возвращает Future
Topic создаётся автоматически на первом publish (если broker настроен auto.create.topics.enable=true).
Partition assignment по key
Если у message есть key, Kafka roughes хеширует и pages в конкретный partition:
# Все messages для user 42 → один partition (preserve order)
await producer.send("user-events", value=event_data, key=b"42")
Это критично когда нужен order для конкретного entity. User events в одном partition - всегда в правильном порядке. Про acks и гарантии доставки - в уроке про delivery semantics.
Без key - round-robin или sticky partitioner.
Consumer
from aiokafka import AIOKafkaConsumer
import json
async def consume():
consumer = AIOKafkaConsumer(
"user-events",
bootstrap_servers="localhost:9092",
group_id="my-group",
auto_offset_reset="earliest", # с начала если нет saved offset
)
await consumer.start()
try:
async for message in consumer:
data = json.loads(message.value.decode())
print(f"Got {data} from partition {message.partition} offset {message.offset}")
finally:
await consumer.stop()
asyncio.run(consume())
group_id- consumer group. Если несколько consumers с тем же group_id, partitions распределятся между нимиauto_offset_reset="earliest"- если у group нет committed offset, начать с начала."latest"- только new messages
Consumer groups - параллельность
Топик с 4 partitions. Consumer group my-group с 2 consumers: A читает P0, P1; B читает P2, P3. Если добавить третий Consumer C - Kafka сделает rebalance: A читает P0, B читает P1 и P2, C читает P3.
Rebalancing автоматически. Один partition → один consumer в group (для preserving order). Если consumers > partitions - некоторые consumers idle. Scaling вверх требует больше partitions.
Multiple consumer groups
Каждая group имеет independent offsets:
# Analytics group - track all events
consumer_analytics = AIOKafkaConsumer("user-events", group_id="analytics", ...)
# Email service - send notifications
consumer_email = AIOKafkaConsumer("user-events", group_id="email-service", ...)
# Search indexer
consumer_search = AIOKafkaConsumer("user-events", group_id="search", ...)
Все 3 groups читают все messages независимо. Это даёт fan-out: один event → много обработок без копирования. Очень мощный паттерн для event-driven архитектуры.
Offset management
consumer = AIOKafkaConsumer(
"user-events",
bootstrap_servers="localhost:9092",
group_id="my-group",
enable_auto_commit=True, # default - commit periodically
auto_commit_interval_ms=5000, # каждые 5 sec
)
С auto-commit: после consume и process, periodic commit offset. Если consumer crash до commit - на restart re-read messages с last commit (возможны duplicates - at-least-once).
Manual commit для precise control:
consumer = AIOKafkaConsumer(
"user-events",
group_id="my-group",
enable_auto_commit=False,
)
await consumer.start()
async for message in consumer:
try:
await process(message.value)
await consumer.commit() # commit после успешной обработки
except Exception:
# не committing - повторим на retry
logger.exception("Failed")
При manual: гарантированно not lose messages, но возможны duplicates.
Idempotent producer
producer = AIOKafkaProducer(
bootstrap_servers="localhost:9092",
enable_idempotence=True, # not send duplicates on retry
acks="all", # ждёт acknowledgement всех replicas
max_in_flight_requests_per_connection=5,
)
Idempotent producer guarantees: если retry due to network, broker не получит duplicate. Это not full exactly-once, но minimum required для надёжной доставки. Стандарт для production producers.
Exactly-once (Kafka transactions)
Combine producer + consumer для read-process-write workflow:
producer = AIOKafkaProducer(
transactional_id="my-transactional-id",
enable_idempotence=True,
)
async with producer.transaction():
# Process input
for input_message in input_consumer:
result = process(input_message)
await producer.send("output-topic", value=result)
# Commit input offsets как part of transaction
await producer.send_offsets_to_transaction(consumed_offsets, group_id)
Это даёт exactly-once в Kafka-to-Kafka workflow. Сложно, требует understanding. Для big data pipelines (Flink, Spark на Kafka) часто используется.
Replay - чтение с начала
consumer = AIOKafkaConsumer(
"events",
bootstrap_servers="localhost:9092",
group_id="new-analytics-team", # new group
auto_offset_reset="earliest", # с самого начала
)
Новый consumer group начинает с beginning топика (auto_offset_reset="earliest"). Можно перечитать всю историю - полезно для:
- Backfill новой analytics системы
- Rebuild state после bug fix
- Onboard новый downstream consumer
Это и есть power Kafka - retention позволяет replay.
Seek - перейти на специфический offset
from aiokafka import TopicPartition
partition = TopicPartition("events", 0)
await consumer.seek(partition, 12345) # offset 12345 в partition 0
# Или к specific время
offsets = await consumer.offsets_for_times({partition: timestamp_ms})
await consumer.seek(partition, offsets[partition].offset)
Можно navigate по historical data. Для debugging, replay specific events, recovery после bug.
Schema registry
Для production со sting evolution schema нужен registry. Avro - стандарт в Kafka ecosystem:
from confluent_kafka.schema_registry import SchemaRegistryClient
from confluent_kafka.schema_registry.avro import AvroSerializer
# Schema (Avro JSON)
schema_str = """
{
"type": "record",
"name": "User",
"fields": [
{"name": "id", "type": "int"},
{"name": "name", "type": "string"}
]
}
"""
sr_client = SchemaRegistryClient({"url": "http://schema-registry:8081"})
serializer = AvroSerializer(sr_client, schema_str, ...)
Schema evolution: добавить optional field, deprecate field - всё managed через registry. Producers и consumers могут быть на разных версиях schema.
aiokafka сам по себе не имеет schema registry support - нужны дополнительные library (confluent-kafka-python имеет).
Использование с FastAPI
from contextlib import asynccontextmanager
from fastapi import FastAPI
from aiokafka import AIOKafkaProducer
@asynccontextmanager
async def lifespan(app: FastAPI):
app.state.producer = AIOKafkaProducer(bootstrap_servers="kafka:9092")
await app.state.producer.start()
yield
await app.state.producer.stop()
app = FastAPI(lifespan=lifespan)
@app.post("/events")
async def create_event(data: EventCreate, producer = Depends(get_producer)):
event = await save_event_to_db(data)
await producer.send(
"events",
value=event.model_dump_json().encode(),
key=str(event.user_id).encode(),
)
return event
Один producer на весь app - thread-safe, можно share.
Standalone consumer (worker)
# consumer_worker.py
import asyncio
import json
import logging
from aiokafka import AIOKafkaConsumer
logger = logging.getLogger(__name__)
async def process_event(data: dict):
# business логика
user_id = data["user_id"]
event_type = data["event"]
# send email, update analytics, etc
async def main():
consumer = AIOKafkaConsumer(
"user-events",
bootstrap_servers="kafka:9092",
group_id="email-service",
auto_offset_reset="earliest",
enable_auto_commit=False,
)
await consumer.start()
try:
async for message in consumer:
try:
data = json.loads(message.value.decode())
await process_event(data)
await consumer.commit()
except Exception:
logger.exception("Failed to process")
# не committed - retry
finally:
await consumer.stop()
if __name__ == "__main__":
asyncio.run(main())
Запускается отдельно от web. Можно scale: запустить N инстансов одного group - Kafka rebalance.
Producing best practices
- Use idempotent producer (enable_idempotence=True) - предотвращает duplicates на retry
- acks=all - wait for all replicas. Slow но gives durability
- Compression (compression_type="gzip" / "snappy" / "lz4") - reduces network and disk
- Batch - aiokafka batches automatically для throughput
- Key для ordering - если порядок важен по entity, use entity_id как key
Consuming best practices
- Manual commit для critical processing (after successful work)
- Idempotent processing (at-least-once)
- Не commit перед work - lose messages если crash
- Handle exceptions - log и either retry or DLQ-style
- Monitor lag (Prometheus Kafka exporter) - индикатор когда отстаёт от tip
Distributed системы considerations
Replication factor - сколько копий каждого partition (для fault tolerance):
- 1 - нет HA (lose broker = lose data)
- 3 - standard для production (tolerate 1 broker failure)
Min in-sync replicas - минимум живых replicas для accepting writes. Защита от data loss.
Retention - как долго хранятся data:
- Time based - 7 дней, 30 дней
- Size based - 100GB на partition
После expiry data удаляется. Replay не возможен старше retention.
Когда Kafka vs RabbitMQ
| Use case | Choose |
|---|---|
| Простой task queue (send email) | RabbitMQ + Celery |
| Event sourcing | Kafka |
| High throughput (>100K msg/s) | Kafka |
| Complex routing rules | RabbitMQ |
| Replay history | Kafka |
| Низкая operational complexity | RabbitMQ или managed Kafka |
| Stream processing (analytics в real-time) | Kafka |
Можно использовать оба: RabbitMQ для tasks, Kafka для events. Не mutually exclusive.
Распространённые ошибки
1. Single partition
# topic с partitions=1
producer.send("events", value=data)
Один partition = один consumer max работает = no scaling. Создавай topics с 10+ partitions для будущего scaling. Можно потом увеличить, но не уменьшить partitions.
2. Commit перед обработкой
async for message in consumer:
await consumer.commit() # COMMIT перед processing
await process(message) # если упадёт - lose
Always commit AFTER successful processing.
3. No idempotence в producer
producer = AIOKafkaProducer(bootstrap_servers="...") # no idempotence
На retry duplicate messages. enable_idempotence=True - стандарт.
4. Один partition для всех related events
Если key = "user-1" для всего, все user events в одном partition. Если processing slow для одного user - вся обработка stuck. Используй key wisely - granularity для parallelism.
5. Игнорировать lag
Lag = разница между producer offset и consumer offset. Если растёт - consumer не успевает. Симптом будущих проблем. Monitor и scale consumers заранее.
Сравнение с Go и PHP
В Go - confluent-kafka-go (wraps librdkafka, fast) или segmentio/kafka-go (pure Go).
writer := kafka.NewWriter(kafka.WriterConfig{Brokers: []string{"localhost:9092"}})
writer.WriteMessages(ctx, kafka.Message{Topic: "events", Value: []byte("hello")})
В PHP - php-rdkafka. Меньше популярен чем Go/Python для streaming.
aiokafka хороший pure-Python вариант. confluent-kafka-python (тоже wraps librdkafka) faster, но более сложная установка (C dependencies).
Мини-задание
- Producer/consumer одной задачей:
import asyncio
import json
from aiokafka import AIOKafkaProducer, AIOKafkaConsumer
async def producer_task():
producer = AIOKafkaProducer(bootstrap_servers="localhost:9092")
await producer.start()
try:
for i in range(5):
await producer.send_and_wait("test", json.dumps({"i": i}).encode())
print(f"Sent {i}")
await asyncio.sleep(1)
finally:
await producer.stop()
async def consumer_task():
consumer = AIOKafkaConsumer(
"test",
bootstrap_servers="localhost:9092",
group_id="test-group",
auto_offset_reset="earliest",
)
await consumer.start()
try:
async for message in consumer:
print(f"Got: {message.value}")
finally:
await consumer.stop()
# Запуск раздельно (в разных терминалах):
# asyncio.run(producer_task())
# asyncio.run(consumer_task())
- Multiple consumers с одним group_id - test partition assignment:
# Создать топик с 3 partitions
kafka-topics --create --topic events --partitions 3 --bootstrap-server localhost:9092
# В разных терминалах запустить 2-3 экземпляра consumer_task()
# Kafka rebalance - каждый получит свой subset partitions
- Replay topic:
async def replay():
consumer = AIOKafkaConsumer(
"test",
bootstrap_servers="localhost:9092",
group_id="replay-group", # new group
auto_offset_reset="earliest",
)
await consumer.start()
count = 0
try:
async for message in consumer:
count += 1
print(f"Replayed: {message.value}")
if count >= 100:
break
finally:
await consumer.stop()
Что дальше
Освоили Kafka. В последнем уроке трека - consumer patterns: идемпотентность, retry с backoff, outbox pattern, graceful shutdown, task queues (Celery, ARQ).