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

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

В предыдущей части мы разобрались с Partition и Consumer Group.

Теперь осталось понять самое интересное:

Что физически хранится внутри Kafka?

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

На самом деле 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.

Почему это настолько важно

Большинство баз данных работает иначе.

Например, UPDATE:

UPDATE users
SET balance = balance - 100

означает:

  • найти страницу на диске;
  • считать её;
  • изменить;
  • записать обратно.

Каждая такая операция требует случайного доступа к диску (Random I/O).

А случайные операции — самые медленные.

Kafka полностью избегает этой проблемы.

Вместо изменения существующих данных она всегда пишет новую запись в конец журнала.

...

Event 100

Event 101

Event 102

Event 103

<- сюда записывается новое сообщение

Благодаря последовательной записи современные SSD и HDD работают значительно эффективнее.

Сообщение в Kafka — это Record

Каждая запись называется Record.

Она состоит из нескольких частей.

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

Рассмотрим каждую подробнее.

Value

Это единственная обязательная часть сообщения.

Именно здесь находится полезная нагрузка.

Например:

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

Kafka совершенно не интересует содержимое Value.

Это может быть:

  • JSON;
  • Avro;
  • Protobuf;
  • XML;
  • обычная строка;
  • бинарные данные.

Для Kafka это просто массив байтов.

Key

Если Value отвечает за содержимое сообщения, то Key определяет, куда сообщение попадёт.

Например:

  • Key = UserID
  • Key = OrderID
  • Key = AccountID

Producer вычисляет хэш ключа и получает номер Partition.

hash(key)
   |
   v
Partition 2

Поэтому все сообщения одного пользователя окажутся в одной партиции.

Именно благодаря Key Kafka сохраняет порядок событий.

Timestamp

Kafka автоматически сохраняет время записи.

Например:

2026-03-15 10:21:54

Это позволяет:

  • строить аналитику;
  • выполнять оконные операции;
  • фильтровать события по времени;
  • реализовывать потоковую обработку.

Headers

Headers — это дополнительные метаданные.

Например:

  • TraceID
  • Content-Type
  • Region
  • RequestID

Очень часто сюда помещают идентификатор трассировки для OpenTelemetry.

Тогда запрос можно проследить через десятки микросервисов.

Теперь самое интересное — Offset

Каждое сообщение получает уникальный номер внутри Partition.

Он называется Offset.

Partition 0

Offset 0

Offset 1

Offset 2

Offset 3

Offset 4

Важно понимать одну вещь.

Offset — это не глобальный идентификатор Kafka.

Он существует только внутри конкретной Partition.

Partition 0

0
1
2
3

--------------

Partition 1

0
1
2
3

Две разные Partition могут одновременно иметь сообщение с Offset = 15.

И это абсолютно нормально.

Зачем вообще нужен Offset?

Представим Consumer.

Он прочитал первые три сообщения.

Offset

0  [x]

1  [x]

2  [x]

3

4

5

Теперь Kafka хранит информацию:

Consumer A

Offset = 3

Это означает:

Следующее сообщение начинается с Offset 3.

Именно поэтому Kafka не удаляет сообщения после чтения.

Она лишь запоминает, до какого места дошёл Consumer.

Что произойдёт после перезапуска?

Допустим, сервер упал.

Через минуту он снова запускается.

Kafka смотрит:

Consumer Offset = 58213

И продолжает чтение именно отсюда.

Никаких потерь.

Никаких повторных чтений, если стратегия подтверждений настроена корректно.

Можно ли прочитать историю заново?

Да.

И это одна из самых мощных возможностей Kafka.

Например:

Consumer Offset

1 500 000

Можно выполнить:

Reset Offset = 0

И Consumer снова прочитает абсолютно всю историю.

0

1

2

3

4

...

1 500 000

Именно благодаря этому работают:

  • Event Sourcing;
  • восстановление аналитики;
  • переиндексация поисковых систем;
  • обучение ML-моделей на исторических данных.

В большинстве классических очередей такое невозможно.

Что происходит при записи сообщения?

Рассмотрим полный путь.

Producer
   |
   v
Topic
   |
   v
Partition
   |
   v
Append Log
   |
   v
Disk

Producer отправляет сообщение.

Kafka вычисляет Partition.

Находит Broker, который является лидером этой Partition.

После этого сообщение просто дописывается в конец файла.

Никаких сложных операций.

Никаких UPDATE.

Никаких DELETE.

Всего одна последовательная запись.

Почему Kafka не использует UUID?

Новички часто задают вопрос:

Почему бы просто не дать каждому сообщению UUID?

Ответ прост.

UUID нужен человеку.

Kafka важна скорость.

Offset обладает сразу несколькими преимуществами:

  • занимает меньше памяти;
  • возрастает последовательно;
  • легко сравнивается;
  • позволяет быстро находить позицию в журнале.

Поэтому Offset значительно эффективнее случайных идентификаторов.

Что дальше?

Теперь мы понимаем, как Kafka хранит сообщения:

  • данные организованы в append-only log;
  • каждое сообщение — это Record;
  • порядок определяется Offset;
  • Consumer хранит только позицию чтения, а не сами сообщения.

Но остаётся ещё один важный вопрос:

Что произойдёт, если Broker выйдет из строя?

В следующей части разберём репликацию, роли Leader и Follower, работу ISR (In-Sync Replicas) и выбор нового лидера при сбое. Проще говоря, посмотрим, почему падение одного сервера не должно превращаться в коллективное чтение логов с холодным потом.