Kafka: Record, Offset и Append-Only Log — Часть 3

Содержание
Коротко: Внутри Partition Kafka хранит append-only log. Offset задаёт позицию записи, а Consumer хранит только свою точку чтения.

В предыдущей части мы разобрались с Partition и Consumer Group. Теперь посмотрим, что Kafka физически хранит внутри.

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

Partition — это “просто” лог

Каждая Partition — упорядоченная последовательность записей:

Partition 0

Offset 0  User Registered
Offset 1  Order Created
Offset 2  Payment Success
Offset 3  Email Sent
Offset 4  Delivery Started

Это одна из причин высокой производительности Kafka. Новые записи добавляются только в конец журнала, поэтому Kafka не приходится искать место для вставки, сдвигать уже записанные данные или постоянно перестраивать файл. Запись выполняется последовательно, крупными блоками, а такой сценарий хорошо использует дисковый кэш операционной системы и заметно быстрее случайных операций чтения и записи. Этот принцип называют append-only log, и он является одной из причин высокой пропускной способности Kafka.

Что представляет собой сообщение

Сообщение в Kafka называется Record. Упрощённо его можно представить так:

+-----------------------+
| Key                   |
| Value                 |
| Timestamp             |
| Headers               |
+-----------------------+

Value — полезная нагрузка сообщения. JSON, Avro, Protobuf, строка или произвольные бинарные данные — для Kafka это просто байты.

Например:

{
  "orderId": 1523,
  "status": "paid",
  "price": 4999
}

Как мы уже говорили ранее – Key часто используется при выборе партиции. Например, ключом может быть UserID или OrderID.

При стандартном key-based partitioning одинаковый ключ будет направляться в одну и ту же партицию, пока количество партиций не меняется. Поэтому события одного объекта можно обрабатывать в правильном порядке.

Timestamp — временная метка записи. Её смысл зависит от настройки топика message.timestamp.type:

  • CreateTime — время, которое задаёт producer. Приложение может явно указать момент события, например время оформления заказа. Если время не указано, producer подставляет время создания записи.
  • LogAppendTime — время добавления записи в журнал Kafka. Его устанавливает брокер, заменяя временную метку, присланную producer.

Например, заказ оформлен в 12:00, а запись о нём попала в Kafka в 12:05. При CreateTime она может сохранить метку 12:00, если приложение передало время заказа. При LogAppendTime метка будет 12:05. Поэтому Timestamp нельзя автоматически считать временем самого события: нужно знать конфигурацию топика и то, как producer заполняет это поле.

Headers позволяют передавать дополнительные метаданные: например, TraceID, Content-Type или RequestID.

Offset

После попадания в партицию запись получает Offset — её позицию внутри журнала. Offset не является глобальным ID сообщения. Он уникален только в пределах конкретной партиции.

Здесь проявляется важное отличие Kafka от традиционной очереди: чтение сообщения не удаляет его из журнала. Kafka хранит записи согласно политике хранения, поэтому их можно прочитать повторно, а разные приложения могут независимо обрабатывать один и тот же поток событий.

alt text

Для Consumer Group Kafka хранит committed offset. Consumer после перезапуска может продолжить чтение с той же позиции, на которой остановился. При этом сами записи остаются в Kafka столько, сколько позволяют настройки retention. А значит, историю можно прочитать повторно. Это удобно, когда нужно заново построить аналитику, переиндексировать данные или повторно обработать события после изменения логики приложения.

Что происходит при записи

Если сильно упростить путь сообщения — Producer отправляет record. Partitioner выбирает партицию, запрос попадает на broker, который является её leader, а запись добавляется в лог.

alt text

Kafka не нужно искать старую запись и изменять её на месте. Основной сценарий — дописать новые данные в конец. И получается вместо постоянного изменения состояния мы получаем последовательную историю событий.

Offset естественным образом отвечает на вопросы:

  • Что было раньше?
  • Что читать следующим?
  • До какого места дошёл consumer?
  • Откуда начать повторное чтение?

Поэтому в Kafka основной адрес записи выглядит примерно так:

Topic + Partition + Offset

Например:

orders / partition 2 / offset 58213

Но пока здесь есть очевидная проблема. Если вся партиция физически находится на одном broker, что произойдёт, когда этот broker упадёт? Для этого в Kafka существует репликация.

В следующей части разберём Leader, Follower, ISR и посмотрим, как Kafka переживает потерю broker.