📨 Apache Kafka

Полный справочник: понятия, топики и партиции, продюсер и консьюмер, группы и offset, гарантии доставки, настройки, CLI, Connect, Streams, схемы, эксплуатация и диагностика.

Шпаргалки · Инфраструктура · #kafka #messaging #streaming #event-driven

Что это

Kafka распределённый журнал событий: сообщения дописываются в конец, хранятся заданное время, читаются много раз многими потребителями. Применение: потоки событий, интеграция сервисов, сбор логов и метрик, CDC (захват изменений БД), event sourcing, потоковая аналитика.

Основные понятия

Термин Смысл
Broker сервер Kafka; несколько брокеров образуют кластер
Topic именованный журнал сообщений
Partition часть топика: упорядоченный неизменяемый лог; порядок гарантирован внутри партиции
Offset порядковый номер сообщения в партиции
Record сообщение: key, value, headers, timestamp
Producer / Consumer отправитель / читатель
Consumer group группа читателей: партиции делятся между ними, каждая партиция читается одним членом
Replication factor число копий партиции на разных брокерах
Leader / follower реплика, обслуживающая запись и чтение / копии
ISR синхронные реплики
Controller брокер, управляющий метаданными (KRaft заменил ZooKeeper)

Ключ сообщения определяет партицию (hash(key) % partitions): одинаковый ключ попадает в одну партицию, значит порядок событий одного ключа сохраняется. Без ключа сообщения распределяются равномерно.

Запуск (KRaft, без ZooKeeper)

KAFKA_CLUSTER_ID="$(bin/kafka-storage.sh random-uuid)"
bin/kafka-storage.sh format -t $KAFKA_CLUSTER_ID -c config/kraft/server.properties
bin/kafka-server-start.sh config/kraft/server.properties

Docker: образ apache/kafka, Compose с KAFKA_PROCESS_ROLES=broker,controller. Порты: 9092 (клиенты), 9093 (контроллер). Параметр advertised.listeners должен быть адресом, доступным клиенту (частая причина «не подключается из Docker»).

Консольные утилиты

BS="--bootstrap-server localhost:9092"
kafka-topics.sh $BS --create --topic orders --partitions 6 --replication-factor 3 --config retention.ms=604800000
kafka-topics.sh $BS --list                       kafka-topics.sh $BS --describe --topic orders
kafka-topics.sh $BS --alter --topic orders --partitions 12        # партиции только увеличивать
kafka-configs.sh $BS --entity-type topics --entity-name orders --alter --add-config retention.ms=86400000
kafka-topics.sh $BS --delete --topic orders

kafka-console-producer.sh $BS --topic orders --property parse.key=true --property key.separator=:
kafka-console-consumer.sh $BS --topic orders --from-beginning --group g1 --property print.key=true --property print.timestamp=true

kafka-consumer-groups.sh $BS --list
kafka-consumer-groups.sh $BS --describe --group g1                # LAG по партициям
kafka-consumer-groups.sh $BS --group g1 --topic orders --reset-offsets --to-earliest --dry-run     # затем --execute
kafka-consumer-groups.sh $BS --group g1 --reset-offsets --to-datetime 2026-01-01T00:00:00.000 --all-topics --execute
kafka-log-dirs.sh $BS --describe                  kafka-get-offsets.sh $BS --topic orders
kafka-producer-perf-test.sh --topic t --num-records 1000000 --record-size 1000 --throughput -1 --producer-props bootstrap.servers=localhost:9092

Сброс offset применяйте при остановленных потребителях группы.

Продюсер

Properties p = new Properties();
p.put("bootstrap.servers", "localhost:9092");
p.put("key.serializer", StringSerializer.class.getName());
p.put("value.serializer", StringSerializer.class.getName());
p.put("acks", "all");
p.put("enable.idempotence", true);
p.put("compression.type", "zstd");
p.put("linger.ms", 10);   p.put("batch.size", 65536);
try (var producer = new KafkaProducer<String, String>(p)) {
    producer.send(new ProducerRecord<>("orders", "order-42", "{...}"), (meta, ex) -> {
        if (ex != null) log.error("не отправлено", ex);
    });
    producer.flush();
}
Настройка Зачем
acks=0 / 1 / all без подтверждения / лидер / все синхронные реплики (надёжно)
enable.idempotence=true без дублей при повторах
retries, delivery.timeout.ms повторные попытки
linger.ms, batch.size накопление пакетов: пропускная способность против задержки
compression.type lz4, zstd, snappy, gzip
max.in.flight.requests.per.connection сохранение порядка при повторах (не больше 5 с идемпотентностью)
transactional.id транзакции

Консьюмер

Properties c = new Properties();
c.put("bootstrap.servers", "localhost:9092");
c.put("group.id", "billing");
c.put("key.deserializer", StringDeserializer.class.getName());
c.put("value.deserializer", StringDeserializer.class.getName());
c.put("auto.offset.reset", "earliest");
c.put("enable.auto.commit", false);
try (var consumer = new KafkaConsumer<String, String>(c)) {
    consumer.subscribe(List.of("orders"));
    while (true) {
        for (var r : consumer.poll(Duration.ofMillis(500))) {
            process(r.key(), r.value());
        }
        consumer.commitSync();           // фиксируем после успешной обработки
    }
}
Настройка Зачем
group.id группа потребителей
auto.offset.reset earliest / latest для новой группы
enable.auto.commit авто-фиксация offset (выключайте для надёжности)
max.poll.records, max.poll.interval.ms размер порции и допустимая пауза между poll
session.timeout.ms, heartbeat.interval.ms обнаружение упавших потребителей
isolation.level=read_committed читать только зафиксированные транзакции
partition.assignment.strategy cooperative-sticky уменьшает остановки при ребалансировке

Гарантии доставки

Гарантия Как получить
At most once фиксировать offset до обработки (возможна потеря)
At least once фиксировать после обработки (возможны дубли): обработка должна быть идемпотентной
Exactly once идемпотентный продюсер + транзакции + read_committed (внутри Kafka), для внешних систем нужна идемпотентность

Надёжная запись: acks=all, min.insync.replicas=2, replication.factor=3, unclean.leader.election.enable=false.

Ребалансировка

Происходит при входе или выходе потребителя, изменении числа партиций. Во время неё чтение приостанавливается. Снижайте: статическое членство (group.instance.id), cooperative-sticky, не блокируйте poll долгой обработкой.

Хранение и очистка

Параметр Смысл
retention.ms / retention.bytes сколько и сколько хранить (по умолчанию 7 дней)
cleanup.policy=delete удалять старые сегменты
cleanup.policy=compact оставлять последнее значение для каждого ключа (журнал состояния)
segment.bytes, segment.ms размер сегмента
min.insync.replicas минимум синхронных реплик для записи с acks=all
max.message.bytes максимальный размер сообщения (по умолчанию ~1 МБ)

Для удаления записи в compact-топике отправьте сообщение с тем же ключом и null значением (tombstone).

Как выбрать число партиций

Партиции задают верхнюю границу параллелизма группы потребителей. Ориентир: целевая пропускная способность / пропускная способность одной партиции, с запасом на рост (6–50 для типичных топиков). Слишком много партиций увеличивают нагрузку на брокеров и время восстановления. Число партиций можно увеличить, но уменьшить нельзя, при увеличении меняется распределение ключей.

Схемы и форматы

JSON, Avro, Protobuf, JSON Schema. Schema Registry (Confluent, Apicurio) хранит схемы и проверяет совместимость (BACKWARD, FORWARD, FULL). Меняйте схемы совместимо: добавляйте поля с значениями по умолчанию.

Экосистема

Компонент Назначение
Kafka Connect готовые коннекторы источников и приёмников (БД, S3, Elasticsearch)
Debezium CDC из PostgreSQL, MySQL, MongoDB в Kafka
Kafka Streams библиотека потоковой обработки в приложении (Java)
ksqlDB, Flink, Spark Structured Streaming SQL и потоковая аналитика
MirrorMaker 2 репликация между кластерами
Kafka UI, AKHQ, Redpanda Console, Conduktor веб-интерфейсы
Клиенты Java, librdkafka (C), Python (confluent-kafka), Go (franz-go, sarama), Node (kafkajs), .NET (Confluent.Kafka), Spring Kafka

Диагностика

Симптом Причина и действия
Растёт LAG потребители не успевают: добавьте партиции и потребителей, ускорьте обработку, проверьте max.poll.records
Постоянные ребалансировки долгая обработка между poll, нестабильная сеть, малые таймауты
Нет подключения advertised.listeners, DNS, файрвол, SASL / TLS настройки
Дубли at-least-once: сделайте обработку идемпотентной
Потери acks=1, автокоммит до обработки, неполные реплики
NOT_ENOUGH_REPLICAS меньше min.insync.replicas живых реплик
RecordTooLargeException увеличьте max.message.bytes и клиентские лимиты или храните большие данные вне Kafka
Порядок нарушен сообщения одного ключа должны идти в одну партицию; повторы без идемпотентности

Мониторинг

Ключевые метрики: consumer lag, under-replicated partitions, offline partitions, ISR shrinks, request latency, disk usage, network I/O, активный контроллер. Инструменты: JMX + Prometheus (kafka_exporter, jmx_exporter), Grafana, Burrow, Cruise Control (балансировка).

Kafka или RabbitMQ

Kafka: поток событий, повторное чтение, большой объём, строгий порядок по ключу. RabbitMQ: очереди задач, гибкая маршрутизация, подтверждения по сообщениям. Подробнее см. шпаргалку RabbitMQ.