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 - параллельность

Kafka consumer group rebalancing: топик с 4 партициями распределяется между 2 или 3 consumers; при добавлении третьего consumer Kafka автоматически перераспределяет партиции

Топик с 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

  1. Use idempotent producer (enable_idempotence=True) - предотвращает duplicates на retry
  2. acks=all - wait for all replicas. Slow но gives durability
  3. Compression (compression_type="gzip" / "snappy" / "lz4") - reduces network and disk
  4. Batch - aiokafka batches automatically для throughput
  5. Key для ordering - если порядок важен по entity, use entity_id как key

Consuming best practices

  1. Manual commit для critical processing (after successful work)
  2. Idempotent processing (at-least-once)
  3. Не commit перед work - lose messages если crash
  4. Handle exceptions - log и either retry or DLQ-style
  5. 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 caseChoose
Простой task queue (send email)RabbitMQ + Celery
Event sourcingKafka
High throughput (>100K msg/s)Kafka
Complex routing rulesRabbitMQ
Replay historyKafka
Низкая operational complexityRabbitMQ или 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).

Мини-задание

  1. 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())
  1. 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
  1. 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).

Зарегистрируйтесь бесплатно, чтобы пройти квиз, решить задание с автопроверкой и вести прогресс.