Половина проблем с Kafka растёт из неверной ментальной модели: к нему относятся как к очереди, а он — распределённый упорядоченный лог. Сообщение не «извлекается и исчезает», оно лежит в партиции по своему offset, пока не истечёт retention, и его читают сколько угодно потребителей независимо друг от друга. Осознание этой разницы объясняет почти всё остальное в Kafka — почему порядок гарантирован не там, где кажется, почему масштабирование чтения устроено иначе, чем в очереди, и почему «повторно обработать поток» — рутинная операция, а не восстановление после катастрофы.
Это первая, вводная статья серии «Kafka: глубокое погружение» — но сама серия целиком держит уровень advanced/deep-dive, рассчитанный на тех, кто уже так или иначе трогал Kafka руками. Если вы только выбираете брокер и хотите обзор без глубины — начните с обзора брокеров и карты messaging, сюда можно будет вернуться позже. Здесь — базовая модель: лог, топик, партиция, offset, producer, consumer, ключ и порядок. Дальше в серии — consumer groups и ребалансировка, репликация и надёжность, хранение и retention, exactly-once, эксплуатация и экосистема. Вся серия опирается на один живой стенд — трёхброкерный KRaft-кластер с примерами на Go (franz-go) и Java, плюс ops-скрипты, — целиком выложенный в публичном репозитории digital-cookbook, каталог kafka/; все числа и распределения ниже — из этого кластера, не из документации.
В статье
- Лог, а не очередь
- Топики и партиции
- Producer и consumer
- Ключ, партиция, порядок: живой прогон
- Репликация — обзор
- Когда Kafka, а когда очередь или NATS
Лог, а не очередь
В классической очереди (RabbitMQ, SQS) сообщение существует, пока его кто-то не подтвердит: после ack брокер его удаляет. Модель заточена под один логический поток обработки — распределить работу между воркерами, а не раздать всем одну и ту же историю событий.
Kafka устроена иначе. Партиция — это упорядоченный, физически неизменяемый (append-only) лог: новые записи дописываются в конец, а уже записанные никогда не переписываются и не удаляются consumer’ом. Позиция записи в этом логе — offset, монотонно растущее целое число, уникальное в пределах партиции. Чтение не меняет лог: consumer просто продвигает свой указатель offset вперёд, а сама запись остаётся на диске до тех пор, пока её не удалит политика retention (время или размер — подробно в статье про хранение и retention), а не факт прочтения.
Отсюда прямое следствие, которое и рвёт шаблон «Kafka — это просто быстрый RabbitMQ»: один и тот же лог может независимо читать сколько угодно consumer’ов (или групп consumer’ов), каждый — со своей позицией. Новый потребитель, подключившийся через полгода, может прочитать всё с начала — если, конечно, retention ещё не подчистил старые записи. Переигрывание потока (replay) — не аварийная процедура, а штатный сценарий: пересчитать агрегат заново, прогнать данные через изменённую бизнес-логику, восстановить состояние сервиса из событий.
Топики и партиции
Топик — логическое имя потока данных (orders, demo-log). Физически топик — это набор партиций, и именно партиция — единица параллелизма и порядка в Kafka, а не топик целиком. Каждая партиция — свой независимый лог со своей последовательностью offset’ов, начинающейся с нуля.
Число партиций топика задаётся при создании и определяет верхнюю границу параллелизма чтения: у группы consumer’ов не может быть активных потребителей партиции больше, чем самих партиций (подробности — в статье про consumer groups). Увеличить число партиций можно на живом топике, а вот уменьшить — нельзя: партиции не «сливаются» обратно, при необходимости топик пересоздаётся заново. Это осознанное ограничение архитектуры, а не недоработка: партиция — это ещё и единица репликации (обзор — ниже), уменьшение переупорядочило бы уже записанные offset’ы, а Kafka как раз обещает их неизменность. Но и увеличение не бесплатно для порядка: раскладка ключей по партициям завязана на их число (murmur2(key) % partitions — подробнее в разделе про ключ ниже), поэтому после расширения топика будущие записи того же ключа могут начать попадать в другую партицию, чем прежде. Уже записанное остаётся на месте, но гарантия «один ключ — один порядок в одной партиции» держится строго при неизменном числе партиций; если порядок по ключу критичен, топик планируют с запасом партиций заранее, а не расширяют на ходу.
Producer и consumer
Producer отправляет запись в топик, а брокер решает, в какую партицию она попадёт: либо явно по ключу записи (об этом — в следующем разделе), либо, если ключа нет, round-robin/sticky-батчингом по доступным партициям. Consumer читает партицию последовательно с произвольной позиции: с начала (earliest), с текущего конца (latest) или с конкретного offset, к которому можно перемотаться вручную (seek).
Ниже — минимальный фрагмент живого стенда: producer пишет пять сообщений на каждый из шести ключей, чередуя ключи по раундам (а не пачками), чтобы видно было именно свойство партиционера, а не побочный эффект батчинга одного ключа подряд.
for round := 0; round < messagesPerKey; round++ {
for _, k := range keys {
value := fmt.Sprintf("%s-msg-%d", k, round)
rec := &kgo.Record{Topic: topic, Key: []byte(k), Value: []byte(value)}
produced, err := cl.ProduceSync(ctx, rec).First()
if err != nil {
log.Fatalf("produce key=%s: %v", k, err)
}
// produced.Partition и produced.Offset — то, что реально
// назначил брокер, а не то, что предполагал клиент.
}
}Тот же сценарий на kafka-clients — синхронный send().get() вместо ProduceSync, но идея та же: партиция и offset становятся известны только из ответа брокера.
for (int round = 0; round < MESSAGES_PER_KEY; round++) {
for (String key : KEYS) {
String value = key + "-msg-" + round;
RecordMetadata meta = producer.send(new ProducerRecord<>(TOPIC, key, value)).get();
results.add(new Sent(key, value, meta.partition(), meta.offset()));
}
}Consumer со стороны чтения устроен зеркально — вычитывает топик с начала и печатает фактическую позицию каждой записи:
cl, err := kgo.NewClient(
kgo.SeedBrokers(seeds...),
kgo.ConsumeTopics(topic),
kgo.ConsumeResetOffset(kgo.NewOffset().AtStart()),
)
// ...
fetches.EachRecord(func(r *kgo.Record) {
out = append(out, recvRecord{
partition: r.Partition,
offset: r.Offset,
key: string(r.Key),
value: string(r.Value),
})
})Коммит позиции чтения (offset), стратегии автокоммита и связанные с ними потери/дубли — тема следующей статьи про consumer groups и ребалансировку; здесь важно только то, что чтение — это движение указателя, а не удаление данных.
Ключ, партиция, порядок: живой прогон
Kafka гарантирует порядок записей строго внутри одной партиции — offset монотонно растёт с шагом 1, и ничего сильнее. Порядок между разными партициями одного топика не гарантирован вообще: если для бизнес-логики важна последовательность событий одной сущности (заказа, пользователя, устройства), эти события обязаны попадать в одну партицию.
За это отвечает ключ записи. Партиционер по умолчанию считает murmur2(key) % partitions — то есть партиция определяется хешем байтов ключа. Это чистая функция: один и тот же ключ при неизменном числе партиций всегда попадёт в одну и ту же партицию, на любом хосте, в любой момент времени. Но у этого правила есть оборотная сторона, которую легко упустить: гарантия работает в одну сторону — одинаковый ключ → всегда одна партиция, но разные ключи вовсе не обязаны разъехаться по разным партициям. Это два разных утверждения, и второе неверно.
На живом стенде это видно без домыслов. Топик demo-log — 3 партиции, RF=3. Producer отправил по 5 сообщений на каждый из 6 ключей (user-1…user-6), итого 30 записей, чередуя ключи по раундам. Вот фактическое распределение — побайтово одинаковое и на Go (franz-go), и на Java (kafka-clients), потому что оба клиента считают ту же формулу (murmur2(key) & 0x7fffffff) % partitions:
partition 0: (пусто)
partition 1: user-4, user-5 (offsets 0..9)
partition 2: user-1, user-2, user-3, user-6 (offsets 0..19)
отправлено=30 / получено=30 (оба клиента)Шесть ключей схлопнулись в две партиции из трёх — партиции 0 не досталось ни одного ключа вообще. Это не баг стенда и не случайность конкретного прогона: это реальная коллизия хеша на конкретном наборе строк, воспроизводимая на любой машине с теми же ключами и тем же числом партиций (значение детерминированное, в отличие от throughput-чисел ниже по серии, которые host-зависимы). Практический вывод: если партиционирование критично для равномерности нагрузки, не полагайтесь на «случайные» строковые ключи — считайте распределение заранее или проектируйте ключи так, чтобы кардинальность была заметно выше числа партиций.
Дополнительно ассерты стенда независимо проверяют обе стороны: у producer’а и у consumer’а один и тот же ключ всегда лежит в одной партиции, а offset внутри каждой партиции идёт подряд без пропусков — и на Go, и на Java эти проверки прошли зелёными.
3 партиции, RF=3"] --> P0["партиция 0
(пусто)"] T --> P1["партиция 1
offset 0..9
user-4, user-5"] T --> P2["партиция 2
offset 0..19
user-1, user-2, user-3, user-6"] K1["ключ user-1"] -->|murmur2 mod 3| P2 K2["ключ user-2"] -->|murmur2 mod 3| P2 K3["ключ user-3"] -->|murmur2 mod 3| P2 K4["ключ user-4"] -->|murmur2 mod 3| P1 K5["ключ user-5"] -->|murmur2 mod 3| P1 K6["ключ user-6"] -->|murmur2 mod 3| P2 style T fill:#f9f3e3,stroke:#8b7355 style P0 fill:#eee8d8,stroke:#8b7355 style P1 fill:#c9e4c5,stroke:#5b8a5e style P2 fill:#c9e4c5,stroke:#5b8a5e
flowchart LR
T["топик demo-log
3 партиции, RF=3"] --> P0["партиция 0
(пусто)"]
T --> P1["партиция 1
offset 0..9
user-4, user-5"]
T --> P2["партиция 2
offset 0..19
user-1, user-2, user-3, user-6"]
K1["ключ user-1"] -->|murmur2 mod 3| P2
K2["ключ user-2"] -->|murmur2 mod 3| P2
K3["ключ user-3"] -->|murmur2 mod 3| P2
K4["ключ user-4"] -->|murmur2 mod 3| P1
K5["ключ user-5"] -->|murmur2 mod 3| P1
K6["ключ user-6"] -->|murmur2 mod 3| P2
style T fill:#f9f3e3,stroke:#8b7355
style P0 fill:#eee8d8,stroke:#8b7355
style P1 fill:#c9e4c5,stroke:#5b8a5e
style P2 fill:#c9e4c5,stroke:#5b8a5e
Репликация — обзор
Партиция — не только единица параллелизма, но и единица репликации: у каждой партиции есть лидер, принимающий все записи и чтения, и фолловеры, копирующие лог с лидера. Если брокер с лидером партиции выходит из строя, лидерство переходит к одному из фолловеров, синхронизированному на момент отказа (входящему в ISR — in-sync replicas).
Это тема, где легко наделать ошибок конфигурацией: за какие гарантии отвечает acks, что такое min.insync.replicas, что происходит при unclean leader election и почему split-brain на уровне метаданных кластера предотвращает не сама репликация, а KRaft-кворум поверх Raft. Подробный разбор — с реальным падением брокера под нагрузкой и проверкой, что ни одна подтверждённая запись не потерялась, — в статье «Репликация и надёжность Kafka». Здесь достаточно знать: репликация в Kafka не опциональная надстройка, а часть модели партиции с первого дня.
Когда Kafka, а когда очередь или NATS
Лог как модель данных — не универсальный инструмент, а конкретный компромисс, и подходит он не для любой задачи с сообщениями между сервисами:
- Kafka уместен, когда нужен высокий устойчивый throughput, поток событий важен сам по себе (а не только как триггер для одноразовой обработки), несколько независимых потребителей должны читать один и тот же поток на разных скоростях, а переигрывание истории — рабочий сценарий, а не исключение.
- Классическая очередь (RabbitMQ) выигрывает там, где важна гибкая маршрутизация (exchange-типы, приоритеты, TTL на сообщение, dead-letter), а поток по сути один — работу нужно распределить между воркерами и подтвердить обработку, а не хранить историю.
- NATS — там, где на первом месте простота эксплуатации и низкая латентность, особенно если нужен встроенный request/reply или лёгкий pub/sub без постоянного хранения (Core NATS), либо облегчённый лог поверх того же кластера (JetStream) без веса полноценной Kafka-инсталляции.
Подробный разбор этого выбора — с конкретными критериями, а не с «Kafka надёжнее» — в статье «RabbitMQ, Kafka, NATS и Redis: как не выбирать messaging layer по моде». А если нужна навигация по всей теме сразу — карта messaging-landscape-map раскладывает все статьи по осям: модель данных, гарантии доставки, эксплуатация, обработка потоков.
Дальше в серии — как Kafka масштабирует чтение через consumer groups и ребалансировку, как устроена репликация и надёжность, и как партиция ведёт себя на диске в статье про retention и log compaction.
Комментарии