сеньорчикОткрыть в Telegram
← вся теориятеория к собесу · Базы и брокеры в проде

Kafka и RabbitMQ в эксплуатации

Брокеры: очередь и журнал

Положил три задачи в очередь и забрал одну потребителем: длина стала 2, забранное сообщение исчезло, и второй потребитель его уже не получит. Потом записал три события в журнал и прочитал их ДВУМЯ независимыми группами: каждая получила все три, а длина журнала осталась 3. Одни и те же данные, две принципиально разные модели доставки.

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

// Формулировки: «чем очередь отличается от журнала?», «потребители не справляются, что делать?», «как не потерять сообщение?»

Две модели доставки

Брокер - это посредник, который принимает сообщения от одних сервисов и отдаёт другим, чтобы те не зависели друг от друга напрямую. Потребитель - программа, которая эти сообщения читает и обрабатывает. Очередь работает как список задач: сообщение кладут, потребитель забирает, подтверждает обработку - и сообщение удаляется. Если потребитель умер, не подтвердив, сообщение возвращается в очередь и достаётся другому. Это модель «работу надо сделать один раз», и она удобна для задач: отправить письмо, посчитать отчёт.

Журнал устроен иначе: сообщения дописываются в конец и лежат положенный срок независимо от того, кто их прочитал. Каждая группа потребителей ведёт СВОЮ отметку о том, докуда дочитала. Мой замер это показывает буквально: две группы прочитали одни и те же три события, ничего друг у друга не забрав, а длина журнала не изменилась. Это модель «событие произошло, кому интересно - читайте», и она удобна, когда одно событие нужно нескольким потребителям.

// Отсюда и разные последствия. В очереди повторная обработка возможна при сбое подтверждения, значит, обработчик обязан быть готов к повтору. В журнале можно перечитать историю заново, просто отмотав отметку назад, - в очереди такого нет, там сообщение уже удалено.

очередь
сообщение удаляется после подтверждения обработки
журнал
сообщения лежат по сроку хранения, каждая группа читает независимо
отметка чтения
докуда группа потребителей дочитала журнал

Подтверждение и повторная обработка

Проверил механику подтверждений на журнале с группой потребителей. Группа прочитала три события - все три сразу оказались в списке неподтверждённых. Подтвердил одно - осталось два. Пока сообщение не подтверждено, оно числится за конкретным потребителем, и если тот умрёт, его работу заберёт другой.

Отсюда главное свойство, которое надо назвать вслух на собесе: брокеры дают доставку «хотя бы один раз». Строгое «ровно один раз» либо стоит очень дорого, либо не существует в той форме, в какой его хотят. Значит, повтор случится, и обработчик обязан быть к нему готов.

// Делают это двумя способами. Либо операция сама по себе безопасна к повтору: «установить статус в оплачено» можно выполнять сколько угодно. Либо у сообщения есть уникальный ключ, и обработчик хранит список уже обработанных ключей. Второй способ надёжнее, но требует хранилища и срока жизни для этих ключей.

подтверждение
сигнал брокеру «обработано, можно забыть»
хотя бы один раз
гарантия доставки: сообщение точно придёт, но может прийти повторно
безопасность к повтору
свойство обработчика: повторная обработка не портит результат

Параллелизм, отставание и срок хранения

В журнале поток режется на партиции, и это единица параллелизма: одну партицию в группе читает ровно один потребитель. Отсюда простая арифметика: партиций 6, потребителей в группе 10 - работают 6, четверо простаивают. Добавлять потребителей сверх числа партиций бесполезно, а увеличить число партиций задним числом можно не всегда безболезненно: порядок сообщений гарантирован только внутри партиции, и перераспределение его нарушает.

Отставание - это разница между тем, что записано, и тем, что прочитано группой. Считается элементарно: приходит 5000 сообщений в секунду, разбирается 3000, значит, отставание растёт на 2000 в секунду. Через минуту это 120 тысяч сообщений, через полчаса - 3,6 миллиона. Поэтому за отставанием следят не по самому числу, а по тому, растёт оно или падает: «миллион и уменьшается» - это хорошая новость, «сто тысяч и растёт» - плохая.

// Срок хранения тоже считается заранее. При 5000 сообщений в секунду и 1,2 килобайта на сообщение хранение суток требует около 494 гигабайт на диске. Это и ограничивает то самое «перечитаем историю заново»: перечитать можно ровно столько, сколько поместилось.

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

Как отвечать: «Потребители не справляются, отставание растёт. Что делаете?»

Сначала смотрю не на само число, а на то, растёт отставание или уже стабилизировалось. Арифметика простая - если приходит пять тысяч в секунду, а разбирается три, отставание растёт на две тысячи в секунду, то есть через полчаса это три с половиной миллиона сообщений. Дальше три рычага. Первый: добавить потребителей, но это работает только до числа партиций - потребители сверх этого числа просто простаивают. Второй: ускорить обработку одного сообщения, обычно это пакетная запись в базу вместо построчной. Третий: увеличить число партиций, но с оговоркой, что порядок гарантирован только внутри партиции и перераспределение его нарушит. И параллельно проверяю срок хранения: если отставание догоняет его, сообщения начнут пропадать непрочитанными, и это уже потеря данных.

Ответ начинается с производной, называет предел масштабирования и заканчивается риском потери. Про упор в число партиций забывают чаще всего.

На чём валятся

  • Добавляют потребителей сверх числа партиций и удивляются, что скорость не выросла.
  • Рассчитывают на доставку ровно один раз и не делают обработчик безопасным к повтору.
  • Смотрят на абсолютное отставание вместо его роста.
  • Забывают про срок хранения: отставание догоняет его, и сообщения пропадают непрочитанными.
  • Считают, что порядок сообщений гарантирован во всём потоке. Он гарантирован внутри партиции.

Проверьте себя

Пять вопросов из банка по этой подтеме. Всего их 6, остальные разбираются в тренажёре.

  1. #dvo_broker_ops1 / 5
    Чем Kafka принципиально отличается от классической очереди вроде RabbitMQ?
    A)Kafka хранит журнал по сроку, и его читают повторно с любого места
    B)Kafka удаляет сообщение после подтверждения, очередь хранит его вечно
    C)Kafka работает только внутри кластера, очередь доступна снаружи
    D)Kafka теряет порядок сообщений, а очередь его держит
    показать ответ и разбор
    +A)Kafka хранит журнал по сроку, и его читают повторно с любого места

    // разбор: Очередь ориентирована на доставку задач: сообщение забрали, подтвердили — оно исчезло, маршрутизация гибкая. Kafka — журнал: записи лежат в партициях до истечения ретенции, каждая группа потребителей держит свою позицию и может перечитать историю с нужного места. Отсюда разные сценарии: очередь для задач и RPC-подобных потоков, журнал для событий, которые нужны нескольким потребителям и для перепроигрывания.

  2. #dvo_broker_ops2 / 5
    Потребители не успевают за потоком. Добавили ещё пять экземпляров — скорость не выросла. Почему?
    A)Группа потребителей не масштабируется больше трёх экземпляров
    B)Новые экземпляры читают с начала журнала и мешают остальным
    C)Параллелизм ограничен числом партиций
    D)Брокер распределяет нагрузку раз в сутки при перебалансировке
    показать ответ и разбор
    +C)Параллелизм ограничен числом партиций

    // разбор: Партиция обрабатывается одним потребителем в группе — это и даёт порядок внутри неё. Если партиций восемь, девятый экземпляр останется без работы. Расширяют число партиций (заранее, с запасом: уменьшать их нельзя, а добавление меняет распределение ключей) либо ускоряют самого потребителя — пачками, параллельной обработкой внутри, выносом тяжёлых операций.

  3. #dvo_broker_ops3 / 5
    Потребитель лежал сутки. После починки часть событий не нашлась, а диски брокера были на пределе. Что случилось?
    A)Брокер удалил неподтверждённые сообщения после таймаута группы
    B)Позиции чтения группы устарели и сбросились в конец журнала
    C)Продюсер перестал писать при заполнении диска и данные потерялись
    D)Ретенция вышла раньше, чем потребитель успел вернуться
    показать ответ и разбор
    +D)Ретенция вышла раньше, чем потребитель успел вернуться

    // разбор: Журнал хранится ограниченное время или до предельного размера. Пока потребитель лежит, старые сегменты уходят по сроку, а его позиция оказывается за пределами доступного — дальше он начинает с самого раннего доступного смещения, и провал уже не восстановить. Отсюда два обязательных сигнала: отставание потребителя и свободное место на брокерах, а ретенцию выбирают из времени, за которое реально чинят потребителя.

  4. #dvo_broker_ops4 / 5
    Брокер сообщений обычно гарантирует доставку «хотя бы один раз». Что это требует от потребителя?
    A)Ничего: брокер сам даёт доставку ровно один раз
    B)Обрабатывать сообщения строго по одному без параллели
    C)Идемпотентной обработки: повтор не задваивает эффект
    D)Подтверждать получение до обработки сообщения
    показать ответ и разбор
    +C)Идемпотентной обработки: повтор не задваивает эффект

    // разбор: Большинство брокеров дают at-least-once: при сбоях/ретраях одно и то же сообщение может прийти повторно (например, потребитель обработал, но не успел подтвердить — брокер пришлёт снова). Значит, потребитель обязан быть идемпотентным: повторная обработка того же сообщения не должна давать двойной эффект (дедуп по ключу сообщения, upsert, проверка «уже сделано»). Рассчитывать на exactly-once по умолчанию нельзя — её либо нет, либо она дорога и с оговорками.

  5. #dvo_broker_ops5 / 5
    Потребитель падает в середине обработки. От чего зависит, потеряется сообщение или обработается дважды?
    A)От размера батча, который потребитель читает за раз
    B)От момента коммита оффсета относительно обработки
    C)От числа партиций в топике брокера
    D)От того, включено ли сжатие сообщений
    показать ответ и разбор
    +B)От момента коммита оффсета относительно обработки

    // разбор: Судьбу сообщения при падении решает момент коммита оффсета (отметки «прочитано до сюда»). Если закоммитить оффсет ДО обработки и упасть — сообщение считается прочитанным и потеряется (at-most-once). Если закоммитить ПОСЛЕ успешной обработки и упасть до коммита — после рестарта оно придёт снова и обработается повторно (at-least-once). Осознанный выбор обычно — коммит после обработки (не терять) плюс идемпотентность (переварить повтор). Это и есть настройка гарантий доставки.

дальше

Теорию прочитали. Навык ставится повторением

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