Kafka Storage & Semantics¶
Kafka как append-only log¶
Фундаментальная идея Kafka:
partition — это append-only log.
То есть запись не вставляется "в середину" и не перестраивает структуру. Она просто дописывается в конец.
Отсюда следуют важные свойства:
- последовательные append-операции очень быстрые;
- offset естественно растет вперед;
- чтение по диапазонам offset удобно и дешево;
- replay становится естественной возможностью.
Что хранится в partition¶
Внутри partition хранится последовательность records:
Каждый record может содержать:
keyvaluetimestampheaders
Важно:
- Kafka не обязана понимать смысл
value; - для Kafka это в первую очередь байты плюс метаданные;
- смысл структуры сообщений — ответственность producer/consumer.
Offset как позиция в логе¶
Offset — это не message id в бизнес-смысле.
Это:
- номер позиции внутри partition;
- способ сказать "читать с этого места";
- очень компактное состояние consumer'а.
Что важно:
- offset монотонно растет внутри partition;
- offset из одной partition нельзя сравнивать с offset другой partition;
- offset может использоваться для replay и seek.
Сегменты лога¶
Физически Kafka не хранит partition как один бесконечный файл.
Лог partition разбивается на сегменты.
Зачем это нужно:
- удобнее удалять старые данные по retention;
- удобнее делать compaction;
- проще управлять файлами на диске;
- можно чистить tail и head без переписывания всего topic.
То есть:
Логически для consumer это один непрерывный offset space, но физически на диске это набор сегментов.
Retention¶
Kafka не хранит данные вечно по умолчанию.
Retention policy определяет:
- как долго хранить данные;
- или какой общий объем хранить;
- после чего старые сегменты можно удалить.
Основные варианты:
- retention по времени;
- retention по размеру;
- compaction policy.
Важно:
- данные удаляются не после чтения consumer'ом;
- данные удаляются по политике хранения topic'а;
- если consumer отстал слишком сильно, retention может обрезать нужные ему старые записи.
Delete-based retention¶
Самый простой вариант:
- Kafka хранит записи ограниченное время или до лимита размера;
- старые сегменты удаляются.
Это удобно для:
- event streams;
- логов;
- телеметрии;
- пайплайнов, где нужен только ограниченный history window.
Пример логики:
Log compaction¶
Compaction — это другая модель хранения.
Идея:
- Kafka старается сохранить последнюю актуальную запись для каждого key;
- старые версии того же key могут быть со временем вычищены.
Это полезно для:
- changelog topics;
- materialized state;
- кэшей и state restoration;
- компактного восстановления "последнего состояния по ключу".
Пример:
После compaction старые версии могут исчезнуть, а последняя останется.
Важные нюансы compaction¶
1. Compaction не мгновенная¶
Kafka не переписывает лог сразу после каждой записи.
Compaction идет фоново.
Значит:
- какое-то время в логе могут лежать и старые, и новые версии одного key;
- consumer может видеть несколько версий подряд.
2. Compaction не меняет порядок¶
Kafka не reorder'ит записи.
Она:
- сохраняет порядок;
- но может удалить часть старых записей из хвоста/сегментов.
3. Offset остается стабильным¶
Даже если запись compacted away, offset как позиция лога логически остается валидной точкой чтения.
Это важно для корректности consumer model.
4. Tombstone¶
Для compacted topics удаление часто выражается как:
- запись с key;
value = null
Это tombstone.
Он говорит:
- "считай ключ удаленным";
- и позже сам tombstone тоже может быть вычищен.
Replay¶
Одна из сильнейших сторон Kafka — возможность replay.
Consumer может:
- перечитать поток с начала;
- начать с конкретного offset;
- читать с конца;
- завести новую consumer group и пройти весь topic заново.
Это полезно, когда:
- починили баг в consumer'е;
- надо пересчитать projection;
- надо восстановить state store;
- надо загрузить новый downstream сервис историей событий.
Kafka здесь сильно отличается от брокеров, где чтение обычно ведет к окончательному удалению сообщения.
Ordering¶
Что Kafka гарантирует¶
- порядок внутри одной partition сохраняется;
- offset внутри partition отражает порядок append.
Что Kafka не гарантирует¶
- глобальный порядок по всему topic;
- порядок между разными partition;
- порядок между разными producer'ами без учета partition strategy.
Как получить порядок по сущности¶
Если нужен порядок по одной сущности, например по orderId, то обычно делают так:
- выбирают
key = orderId; - гарантируют, что все события одного order идут в одну partition.
Тогда можно ожидать:
OrderCreatedOrderPaidOrderShipped
в правильном порядке внутри этой partition.
Но надо помнить:
- если поменять число partition, распределение key может измениться;
- если producer пишет без key, такого порядка не будет;
- глобального порядка между разными order все равно нет.
Почему Kafka не классическая очередь¶
Kafka похожа на очередь только на очень грубом уровне: "кто-то пишет, кто-то читает".
Но концептуально отличия большие:
В классической очереди¶
- сообщение часто забирается одним consumer;
- после ack сообщение удаляется;
- replay неудобен или отсутствует;
- приоритет на task dispatch.
В Kafka¶
- данные лежат в логе по retention;
- несколько groups могут читать один и тот же topic;
- offset хранится отдельно;
- replay — нормальная функция;
- приоритет на event stream, throughput и durable log.
Начало чтения новой consumer group¶
Когда новая group впервые приходит в topic, важен вопрос:
- читать с начала;
- или только новые записи.
Это связано с auto.offset.reset.
Типичная логика:
earliest— читать максимально с начала доступного retention window;latest— читать только новые поступающие записи.
Практический смысл:
- аналитике и реплейным сервисам часто нужен
earliest; - онлайн-подписчику иногда нужен
latest.
Exactly-once: что здесь важно понимать со стороны хранения¶
Kafka storage model очень хорошо помогает сделать надежный pipeline, потому что:
- записи лежат в durable log;
- offsets компактны;
- commit offset можно хранить в Kafka;
- replay естественно поддерживается.
Но storage model сама по себе не означает:
- что внешний side effect станет exactly-once;
- что consumer никогда не увидит дубль;
- что любое приложение без усилий станет консистентным.
Storage дает хорошие primitives. Правильное приложение должно ими правильно воспользоваться.
Что запомнить¶
- Partition — это append-only log.
- Offset — это позиция в этом логе.
- Kafka хранит данные по retention policy, а не "до первого чтения".
- Compaction полезна для состояния по ключу, а не для обычного event history.
- Replay — одна из ключевых сильных сторон Kafka.
- Порядок гарантирован только внутри partition.