Kafka: как масштабируются Partitions, Consumer Groups и порядок сообщений — Часть 2

Содержание
Коротко: Partitions дают Kafka параллелизм, а Consumer Groups распределяют чтение между обработчиками. Порядок сообщений сохраняется только внутри одной Partition.

В первой части мы познакомились с основными компонентами Kafka:

  • Broker
  • Topic
  • Producer
  • Consumer

Но главный вопрос остался открытым, и он обычно появляется сразу после первой фразы «Kafka выдерживает миллионы сообщений».

Как Kafka способна обрабатывать миллионы сообщений в секунду и при этом не терять порядок данных?

Ответ состоит всего из двух слов:

Partitions и Consumer Groups.

Именно они делают Kafka тем инструментом, которым она является сегодня. Без них Kafka была бы просто очень уверенным в себе лог-файлом.

Почему одного сервера недостаточно

Представим обычный интернет-магазин.

Каждую секунду происходят события:

  • оформление заказа;
  • оплата;
  • отмена;
  • возврат;
  • доставка.

Поначалу всё выглядит так:

Producer
   |
   v
+-------------------+
|       Kafka       |
+-------------------+
   |
   v
Consumer

Система работает отлично. До того момента, когда «немного трафика» внезапно оказывается полноценной нагрузкой.

Но проходит время.

Теперь:

  • 50 000 заказов в минуту;
  • тысячи пользователей одновременно;
  • десятки микросервисов.

Один сервер начинает упираться в:

  • CPU;
  • диск;
  • сеть.

Вертикально масштабироваться можно не бесконечно.

Поэтому Kafka масштабируется горизонтально.

Первая мысль — добавить ещё серверов

Допустим, мы добавили три брокера:

Broker 1
Broker 2
Broker 3

Возникает вопрос.

Как распределять сообщения?

Если случайным образом:

Order #101 -> Broker 2
Order #102 -> Broker 1
Order #103 -> Broker 3

Получается хаос. Не драматический, с молниями и музыкой, а обычный продовый: всё формально работает, но порядок событий уже живёт своей жизнью.

Почему случайное распределение ломает систему

Представим банковский перевод.

Последовательность событий должна быть такой:

Создан счёт
   |
   v
Пополнение
   |
   v
Списание
   |
   v
Закрытие

Но если сообщения попадут на разные серверы и будут обработаны с разной скоростью, Consumer может увидеть:

Закрытие
   |
   v
Пополнение
   |
   v
Создание счёта

Получится неконсистентное состояние.

Поэтому Kafka должна обеспечить два свойства одновременно:

  • высокий параллелизм;
  • сохранение порядка.

Именно для этого существуют Partition.

Что такое Partition

Topic — это всего лишь логическая сущность.

Физически Topic разбивается на несколько независимых журналов.

Например:

Topic: orders

+-------------+
| Partition 0 |
+-------------+

+-------------+
| Partition 1 |
+-------------+

+-------------+
| Partition 2 |
+-------------+

Каждая партиция — это отдельный append-only лог.

Именно туда последовательно записываются сообщения.

Почему Partition делает Kafka быстрее

Теперь вместо одного журнала появляются три.

Partition 0
1
2
3
4
5

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

Partition 1
1
2
3
4

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

Partition 2
1
2
3

Все они могут:

  • записываться одновременно;
  • читаться одновременно;
  • храниться на разных серверах.

Получается почти линейное масштабирование.

Очень важное правило Kafka

Kafka гарантирует порядок сообщений только внутри одной Partition.

Именно поэтому нельзя говорить:

Kafka гарантирует порядок.

Правильнее сказать:

Kafka гарантирует порядок в пределах партиции.

Это один из самых популярных вопросов на собеседованиях.

Как Kafka понимает, куда отправить сообщение

Предположим, есть Topic с четырьмя Partition.

orders

P0
P1
P2
P3

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

Как выбрать партицию?

Ответ зависит от того, есть ли у сообщения Key.

Если Key отсутствует

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

Например:

msg1 -> P0
msg2 -> P1
msg3 -> P2
msg4 -> P3
msg5 -> P0

Нагрузка получается равномерной.

Но порядок связанных сообщений теряется.

Если есть Key

Это самый распространённый сценарий.

Например:

  • OrderID = 52341
  • UserID = 1827
  • AccountID = 9912

Kafka вычисляет хэш ключа.

Упрощённо:

hash(key)
   |
   v
mod
   |
   v
Количество Partition

Например:

hash(52341)
   |
   v
458271
   |
   v
458271 % 3 = 0

Сообщение всегда попадёт в Partition 0.

И завтра.

И через месяц.

И через год.

Это означает, что все события одного заказа всегда будут храниться в одной партиции, а значит — сохранят правильный порядок.

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

Допустим, заказ проходит несколько этапов.

Order Created
   |
   v
Payment Started
   |
   v
Payment Success
   |
   v
Delivery Started
   |
   v
Delivered

Все эти события имеют одинаковый OrderID.

Следовательно:

hash(OrderID)
   |
   v
Partition 2

Consumer увидит события именно в той последовательности, в которой они были записаны.

А если партиций станет больше?

Это один из подводных камней Kafka.

Предположим, было:

3 Partition

Стало:

6 Partition

Формула изменилась.

Теперь:

hash(key) % 6

Некоторые ключи начнут попадать в другие Partition.

Из-за этого при увеличении числа партиций порядок старых и новых сообщений для одного ключа может нарушиться. Поэтому количество Partition обычно стараются планировать заранее.

Теперь поговорим о Consumer

До сих пор у нас был один Consumer.

Partition
   |
   v
Consumer

Но если сообщений становится миллион?

Один Consumer уже не справится.

Consumer Group

Kafka объединяет несколько Consumer в одну группу.

Consumer Group

+------------+
| Consumer A |
+------------+

+------------+
| Consumer B |
+------------+

+------------+
| Consumer C |
+------------+

Теперь Kafka распределяет партиции между ними.

Кто читает какую партицию?

Например:

Partition 0
   |
   v
Consumer A

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

Partition 1
   |
   v
Consumer B

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

Partition 2
   |
   v
Consumer C

Каждая партиция принадлежит только одному Consumer внутри группы.

Благодаря этому сообщение никогда не обработается дважды.

Что если Consumer упал?

Допустим:

Consumer B
    X

Kafka замечает это.

И выполняет rebalance.

Partition 1
   |
   v
Consumer A

или:

Partition 1
   |
   v
Consumer C

Работа продолжается практически автоматически.

Именно поэтому Consumer Group обеспечивает отказоустойчивость без участия разработчика.

Можно ли сделать Consumer больше, чем Partition?

Например:

3 Partition
5 Consumer

Ответ:

Можно.

Но два Consumer будут простаивать.

Kafka никогда не отдаёт одну Partition двум Consumer внутри одной группы.

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

А если нужно читать один Topic разными сервисами?

Очень часто начинающие думают:

Если один Consumer прочитал сообщение, второй его уже не увидит.

Это неверно.

Пусть есть Topic:

orders

И три независимых сервиса:

  • Email;
  • Analytics;
  • Fraud Detection.

Каждый создаёт собственную Consumer Group.

Topic: orders
      |
      v
----------------
Email Group
----------------
Analytics Group
----------------
Fraud Group
----------------

Каждая группа хранит свои собственные Offset и читает Topic независимо от остальных.

В результате одно и то же событие может одновременно:

  • отправить письмо;
  • обновить отчёт;
  • проверить заказ на мошенничество.

И всё это без дублирования сообщений и без участия Producer.

Что дальше?

Теперь мы понимаем две ключевые идеи Kafka:

  • Partition отвечает за масштабирование и порядок сообщений.
  • Consumer Group отвечает за параллельную обработку и отказоустойчивость.

Но остаётся главный вопрос:

Что именно хранится внутри Partition?

В следующей части заглянем внутрь Partition: что такое Record, зачем нужен Offset, как работает append-only log и почему Kafka спокойно пишет на диск, пока многие системы от одной мысли об этом начинают тяжело дышать.