Kafka и RabbitMQ в эксплуатации
Положил три задачи в очередь и забрал одну потребителем: длина стала 2, забранное сообщение исчезло, и второй потребитель его уже не получит. Потом записал три события в журнал и прочитал их ДВУМЯ независимыми группами: каждая получила все три, а длина журнала осталась 3. Одни и те же данные, две принципиально разные модели доставки.
Стержень: очередь удаляет сообщение после обработки, журнал хранит его по сроку и позволяет читать многим независимо; параллелизм у журнала ограничен числом партиций.
// Формулировки: «чем очередь отличается от журнала?», «потребители не справляются, что делать?», «как не потерять сообщение?»
Две модели доставки
Брокер - это посредник, который принимает сообщения от одних сервисов и отдаёт другим, чтобы те не зависели друг от друга напрямую. Потребитель - программа, которая эти сообщения читает и обрабатывает. Очередь работает как список задач: сообщение кладут, потребитель забирает, подтверждает обработку - и сообщение удаляется. Если потребитель умер, не подтвердив, сообщение возвращается в очередь и достаётся другому. Это модель «работу надо сделать один раз», и она удобна для задач: отправить письмо, посчитать отчёт.
Журнал устроен иначе: сообщения дописываются в конец и лежат положенный срок независимо от того, кто их прочитал. Каждая группа потребителей ведёт СВОЮ отметку о том, докуда дочитала. Мой замер это показывает буквально: две группы прочитали одни и те же три события, ничего друг у друга не забрав, а длина журнала не изменилась. Это модель «событие произошло, кому интересно - читайте», и она удобна, когда одно событие нужно нескольким потребителям.
// Отсюда и разные последствия. В очереди повторная обработка возможна при сбое подтверждения, значит, обработчик обязан быть готов к повтору. В журнале можно перечитать историю заново, просто отмотав отметку назад, - в очереди такого нет, там сообщение уже удалено.
- очередь
- сообщение удаляется после подтверждения обработки
- журнал
- сообщения лежат по сроку хранения, каждая группа читает независимо
- отметка чтения
- докуда группа потребителей дочитала журнал
Подтверждение и повторная обработка
Проверил механику подтверждений на журнале с группой потребителей. Группа прочитала три события - все три сразу оказались в списке неподтверждённых. Подтвердил одно - осталось два. Пока сообщение не подтверждено, оно числится за конкретным потребителем, и если тот умрёт, его работу заберёт другой.
Отсюда главное свойство, которое надо назвать вслух на собесе: брокеры дают доставку «хотя бы один раз». Строгое «ровно один раз» либо стоит очень дорого, либо не существует в той форме, в какой его хотят. Значит, повтор случится, и обработчик обязан быть к нему готов.
// Делают это двумя способами. Либо операция сама по себе безопасна к повтору: «установить статус в оплачено» можно выполнять сколько угодно. Либо у сообщения есть уникальный ключ, и обработчик хранит список уже обработанных ключей. Второй способ надёжнее, но требует хранилища и срока жизни для этих ключей.
- подтверждение
- сигнал брокеру «обработано, можно забыть»
- хотя бы один раз
- гарантия доставки: сообщение точно придёт, но может прийти повторно
- безопасность к повтору
- свойство обработчика: повторная обработка не портит результат
Параллелизм, отставание и срок хранения
В журнале поток режется на партиции, и это единица параллелизма: одну партицию в группе читает ровно один потребитель. Отсюда простая арифметика: партиций 6, потребителей в группе 10 - работают 6, четверо простаивают. Добавлять потребителей сверх числа партиций бесполезно, а увеличить число партиций задним числом можно не всегда безболезненно: порядок сообщений гарантирован только внутри партиции, и перераспределение его нарушает.
Отставание - это разница между тем, что записано, и тем, что прочитано группой. Считается элементарно: приходит 5000 сообщений в секунду, разбирается 3000, значит, отставание растёт на 2000 в секунду. Через минуту это 120 тысяч сообщений, через полчаса - 3,6 миллиона. Поэтому за отставанием следят не по самому числу, а по тому, растёт оно или падает: «миллион и уменьшается» - это хорошая новость, «сто тысяч и растёт» - плохая.
// Срок хранения тоже считается заранее. При 5000 сообщений в секунду и 1,2 килобайта на сообщение хранение суток требует около 494 гигабайт на диске. Это и ограничивает то самое «перечитаем историю заново»: перечитать можно ровно столько, сколько поместилось.
- партиция
- часть потока; одну партицию в группе читает ровно один потребитель
- отставание
- разница между записанным и прочитанным; смотреть надо на его рост
- срок хранения
- сколько журнал держит сообщения; определяет глубину перечитывания
Как отвечать: «Потребители не справляются, отставание растёт. Что делаете?»
Сначала смотрю не на само число, а на то, растёт отставание или уже стабилизировалось. Арифметика простая - если приходит пять тысяч в секунду, а разбирается три, отставание растёт на две тысячи в секунду, то есть через полчаса это три с половиной миллиона сообщений. Дальше три рычага. Первый: добавить потребителей, но это работает только до числа партиций - потребители сверх этого числа просто простаивают. Второй: ускорить обработку одного сообщения, обычно это пакетная запись в базу вместо построчной. Третий: увеличить число партиций, но с оговоркой, что порядок гарантирован только внутри партиции и перераспределение его нарушит. И параллельно проверяю срок хранения: если отставание догоняет его, сообщения начнут пропадать непрочитанными, и это уже потеря данных.
Ответ начинается с производной, называет предел масштабирования и заканчивается риском потери. Про упор в число партиций забывают чаще всего.
На чём валятся
- −Добавляют потребителей сверх числа партиций и удивляются, что скорость не выросла.
- −Рассчитывают на доставку ровно один раз и не делают обработчик безопасным к повтору.
- −Смотрят на абсолютное отставание вместо его роста.
- −Забывают про срок хранения: отставание догоняет его, и сообщения пропадают непрочитанными.
- −Считают, что порядок сообщений гарантирован во всём потоке. Он гарантирован внутри партиции.
Проверьте себя
Пять вопросов из банка по этой подтеме. Всего их 6, остальные разбираются в тренажёре.
- Чем Kafka принципиально отличается от классической очереди вроде RabbitMQ?A)Kafka хранит журнал по сроку, и его читают повторно с любого местаB)Kafka удаляет сообщение после подтверждения, очередь хранит его вечноC)Kafka работает только внутри кластера, очередь доступна снаружиD)Kafka теряет порядок сообщений, а очередь его держит
показать ответ и разбор
+A)Kafka хранит журнал по сроку, и его читают повторно с любого места// разбор: Очередь ориентирована на доставку задач: сообщение забрали, подтвердили — оно исчезло, маршрутизация гибкая. Kafka — журнал: записи лежат в партициях до истечения ретенции, каждая группа потребителей держит свою позицию и может перечитать историю с нужного места. Отсюда разные сценарии: очередь для задач и RPC-подобных потоков, журнал для событий, которые нужны нескольким потребителям и для перепроигрывания.
- Потребители не успевают за потоком. Добавили ещё пять экземпляров — скорость не выросла. Почему?A)Группа потребителей не масштабируется больше трёх экземпляровB)Новые экземпляры читают с начала журнала и мешают остальнымC)Параллелизм ограничен числом партицийD)Брокер распределяет нагрузку раз в сутки при перебалансировке
показать ответ и разбор
+C)Параллелизм ограничен числом партиций// разбор: Партиция обрабатывается одним потребителем в группе — это и даёт порядок внутри неё. Если партиций восемь, девятый экземпляр останется без работы. Расширяют число партиций (заранее, с запасом: уменьшать их нельзя, а добавление меняет распределение ключей) либо ускоряют самого потребителя — пачками, параллельной обработкой внутри, выносом тяжёлых операций.
- Потребитель лежал сутки. После починки часть событий не нашлась, а диски брокера были на пределе. Что случилось?A)Брокер удалил неподтверждённые сообщения после таймаута группыB)Позиции чтения группы устарели и сбросились в конец журналаC)Продюсер перестал писать при заполнении диска и данные потерялисьD)Ретенция вышла раньше, чем потребитель успел вернуться
показать ответ и разбор
+D)Ретенция вышла раньше, чем потребитель успел вернуться// разбор: Журнал хранится ограниченное время или до предельного размера. Пока потребитель лежит, старые сегменты уходят по сроку, а его позиция оказывается за пределами доступного — дальше он начинает с самого раннего доступного смещения, и провал уже не восстановить. Отсюда два обязательных сигнала: отставание потребителя и свободное место на брокерах, а ретенцию выбирают из времени, за которое реально чинят потребителя.
- Брокер сообщений обычно гарантирует доставку «хотя бы один раз». Что это требует от потребителя?A)Ничего: брокер сам даёт доставку ровно один разB)Обрабатывать сообщения строго по одному без параллелиC)Идемпотентной обработки: повтор не задваивает эффектD)Подтверждать получение до обработки сообщения
показать ответ и разбор
+C)Идемпотентной обработки: повтор не задваивает эффект// разбор: Большинство брокеров дают at-least-once: при сбоях/ретраях одно и то же сообщение может прийти повторно (например, потребитель обработал, но не успел подтвердить — брокер пришлёт снова). Значит, потребитель обязан быть идемпотентным: повторная обработка того же сообщения не должна давать двойной эффект (дедуп по ключу сообщения, upsert, проверка «уже сделано»). Рассчитывать на exactly-once по умолчанию нельзя — её либо нет, либо она дорога и с оговорками.
- Потребитель падает в середине обработки. От чего зависит, потеряется сообщение или обработается дважды?A)От размера батча, который потребитель читает за разB)От момента коммита оффсета относительно обработкиC)От числа партиций в топике брокераD)От того, включено ли сжатие сообщений
показать ответ и разбор
+B)От момента коммита оффсета относительно обработки// разбор: Судьбу сообщения при падении решает момент коммита оффсета (отметки «прочитано до сюда»). Если закоммитить оффсет ДО обработки и упасть — сообщение считается прочитанным и потеряется (at-most-once). Если закоммитить ПОСЛЕ успешной обработки и упасть до коммита — после рестарта оно придёт снова и обработается повторно (at-least-once). Осознанный выбор обычно — коммит после обработки (не терять) плюс идемпотентность (переварить повтор). Это и есть настройка гарантий доставки.
дальше
Теорию прочитали. Навык ставится повторением
В Сеньорчике эта подтема идёт в ежедневных сессиях: движок возвращает её, пока ответы не станут уверенными, и ведёт прогресс отдельно по каждой подтеме. Теория внутри тоже бесплатна, лимит только на количество вопросов в день.