📨 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.