Очереди и Kafka в Java
Прямой вызов требует, чтобы сосед был жив прямо сейчас. Сообщение через брокера - отдельную службу-посредника, которая принимает сообщения и хранит их до получателя, - этого не требует: отправитель кладёт сообщение и идёт дальше, получатель забирает когда сможет. Всплеск нагрузки сглаживается буфером, а падение получателя перестаёт быть падением отправителя.
Цена - в том, что синхронная ясность пропадает. Появляются дубликаты, порядок перестаёт быть само собой разумеющимся, а ошибка обработки всплывает не там, где её породили.
// Формулировки: «зачем брокер?», «что такое at-least-once?», «как обеспечить порядок?», «куда девать сообщения, которые не обрабатываются?».
Очередь или лог
Брокеры делятся на два семейства, и это первое, что стоит различать. Классическая очередь (RabbitMQ и подобные) отдаёт сообщение потребителю и удаляет его после подтверждения: сообщение существует, пока его не забрали. Лог (Kafka и подобные) хранит поток записей на диске заданное время, а каждый потребитель держит свою позицию чтения.
Отсюда разные возможности. У лога можно перечитать историю заново - поднять новый сервис и дать ему прожевать прошлый месяц событий. У очереди перечитать нечего: обработанное удалено.
// Выбор по задаче. Очередь удобна для заданий, у которых есть исполнитель: разослать письма, сгенерировать отчёт. Лог удобен для событий, у которых потребителей заранее неизвестно сколько: «заказ оплачен» интересен и складу, и аналитике, и уведомлениям, причём каждый читает в своём темпе.
- очередь / лог
- сообщение удаляется после подтверждения / хранится, у каждого читателя своя позиция
Доставка минимум один раз и порядок
Гарантия по умолчанию почти везде - минимум один раз: сообщение точно дойдёт, но возможно не в единственном экземпляре. Дубликаты появляются буднично: отправитель не дождался подтверждения и повторил; получатель обработал сообщение и умер, не успев отметить; сработала перебалансировка потребителей.
Отсюда главное требование к обработчику: он обязан быть идемпотентным, то есть повторная обработка того же сообщения не должна менять итог. Делают это ключом: у сообщения есть идентификатор, обработчик хранит уже обработанные и повтор пропускает, либо операция пишется так, что повтор безвреден по своей природе - установка значения вместо прибавления.
// Теперь про порядок. Поток сообщений режут на разделы - независимые части, которые читаются параллельно разными обработчиками, ради скорости. Порядок гарантируется только ВНУТРИ одного раздела. Поэтому, если события по одному заказу должны идти строго по очереди, номер заказа берут ключом разбиения: тогда все они попадают в один раздел и читаются последовательно. Ждать глобального порядка по всему потоку не стоит: он несовместим с параллельной обработкой.
- минимум один раз
- дойдёт точно, возможно не единожды - обработчик обязан терпеть повторы
- ключ разбиения
- по нему сообщения попадают в один раздел и сохраняют порядок
Ошибки и содержание события
Сообщение, которое не удаётся обработать, нельзя ни терять, ни пытаться обработать вечно. Обычная схема лесенкой: несколько повторов с растущей паузой, и если не вышло - отправка в отдельную очередь для неразобранных, с оповещением дежурному. Без такой очереди одно сломанное сообщение встаёт пробкой и останавливает обработку всех следующих.
Отдельно про содержание. Событие должно нести факт и его данные - «заказ 42 оплачен, сумма такая-то», а не команду «отправь письмо». Команда прибивает отправителя к конкретному получателю; факт позволяет подписаться на него кому угодно и не трогать отправителя.
// И про схему. Формат сообщения - такой же публичный контракт, как и API: у него есть версия, добавлять поля можно, удалять и переименовывать нельзя. В логе это особенно жёстко, потому что старые сообщения лежат месяцами и их будет читать новый код.
- очередь неразобранных
- куда уходит сообщение после исчерпания повторов, с оповещением
Как отвечать: «Доставка минимум один раз: что она требует от потребителя?»
Требует идемпотентности обработчика: повторная обработка того же сообщения не должна менять результат. Дубликаты возникают штатно - отправитель не дождался подтверждения и повторил, потребитель обработал и упал до отметки о прочтении, случилась перебалансировка. Делаю это двумя способами. Первый - ключ: у сообщения есть уникальный идентификатор, я храню уже обработанные и повтор пропускаю. Второй - формулировать операцию так, чтобы повтор был безвреден сам по себе: установить статус в «оплачен» можно сколько угодно раз, а вот прибавить сумму к балансу - нет. Ещё слежу за порядком: он сохраняется только внутри раздела, поэтому события одной сущности отправляю с ключом разбиения по её идентификатору. И обязательно завожу очередь для неразобранных сообщений с оповещением - иначе одно сломанное сообщение встанет пробкой и остановит всю обработку.
Почему это сильный ответ: названы конкретные причины дубликатов, два разных способа добиться идемпотентности и смежные обязательные вещи - порядок по ключу и очередь неразобранных.
На чём валят
- −Считать доставку ровно однократной. По умолчанию она «минимум один раз», и обработчик обязан терпеть повторы.
- −Ждать глобального порядка сообщений. Порядок есть только внутри раздела и только по одному ключу.
- −Работать без очереди неразобранных. Одно сломанное сообщение встаёт пробкой перед всеми следующими.
- −Класть в событие команду вместо факта. Отправитель оказывается прибит к конкретному получателю.
- −Менять формат сообщения как внутреннюю структуру. Это публичный контракт, а в логе старые сообщения живут месяцами.
Проверьте себя
Пять вопросов из банка по этой подтеме. Всего их 17, остальные разбираются в тренажёре.
- Как Kafka обеспечивает порядок сообщений?A)Kafka обеспечивает строгий глобальный порядок всех сообщений по всему топику независимо от партицийB)Порядок в Kafka не сохраняется вообще: потребитель получает сообщения в случайном порядкеC)Порядок гарантирован в пределах ОДНОЙ партиции, а не по всему топикуD)Порядок задаётся временем отправки: брокер сортирует все сообщения по timestamp перед выдачей
показать ответ и разбор
+C)Порядок гарантирован в пределах ОДНОЙ партиции, а не по всему топику// разбор: Kafka-топик делится на ПАРТИЦИИ для параллелизма. Порядок гарантирован ТОЛЬКО В ПРЕДЕЛАХ одной партиции (сообщения читаются по возрастанию offset), а МЕЖДУ партициями порядка нет. Чтобы связанные сообщения шли по порядку (все события одного заказа/пользователя), им задают ОДИН КЛЮЧ — Kafka хешем ключа кладёт их в одну партицию, сохраняя их взаимный порядок. Плата за строгий порядок — меньше параллелизма (одна партиция = один потребитель в группе). Поэтому ключ партиционирования выбирают так, чтобы совместить нужную упорядоченность с равномерным распределением нагрузки.
- Что такое consumer group в Kafka?A)Набор продюсеров, совместно пишущих в один топик, чтобы ускорить отправку сообщений в негоB)Группа топиков, объединённых под одним именем, которые потребитель читает как единый потокC)Резервная копия потребителя: второй экземпляр обрабатывает те же сообщения параллельно для надёжностиD)Группа потребителей, делящих партиции топика для параллельной обработки без дублей внутри группы
показать ответ и разбор
+D)Группа потребителей, делящих партиции топика для параллельной обработки без дублей внутри группы// разбор: Consumer group — механизм МАСШТАБИРОВАНИЯ потребления: партиции топика РАСПРЕДЕЛЯЮТСЯ между участниками группы, каждая партиция закреплена за ОДНИМ потребителем группы — так сообщения обрабатываются параллельно и БЕЗ дублей ВНУТРИ группы. При добавлении/падении потребителя происходит ребаланс (перераспределение партиций). РАЗНЫЕ группы читают один топик НЕЗАВИСИМО (каждая получает все сообщения — это pub/sub-семантика между группами). Число полезных потребителей в группе ограничено числом партиций (лишние простаивают). Прогресс группы хранится в offset'ах.
- Что должно произойти с сообщением, которое раз за разом падает при обработке?A)После N попыток уйти в dead-letter queue, не блокируя остальные сообщенияB)Оставаться в голове очереди и обрабатываться повторно бесконечно, пока однажды всё-таки не пройдётC)Быть молча отброшенным сразу после первой же ошибки, чтобы не задерживать поток обработкиD)Останавливать всю очередь целиком до тех пор, пока администратор не разберёт проблему вручную
показать ответ и разбор
+A)После N попыток уйти в dead-letter queue, не блокируя остальные сообщения// разбор: «Ядовитое» (poison) сообщение — то, что потребитель не может обработать (битый формат, ссылка на отсутствующие данные, баг). Если ретраить его бесконечно в голове очереди — заблокируется вся дальнейшая обработка; если молча выбросить — потеряется факт проблемы. Правильно: после заданного числа неудачных попыток переместить сообщение в отдельную DEAD-LETTER QUEUE (DLQ), где его позже разберут/переобработают, а основной поток продолжит работу. DLQ мониторят (её рост — сигнал проблемы). Это стандартный приём устойчивой асинхронной обработки в RabbitMQ/Kafka (через retry-топики и DLQ).
- Чем модель pub/sub отличается от очереди (точка-точка)?A)Pub/sub доставляет сообщение только одному случайному подписчику, а очередь — сразу всем сразуB)Pub/sub доставляет событие ВСЕМ подписчикам; очередь — одному из потребителейC)Это одно и то же; разница лишь в названии у разных брокеров сообщений без функциональных отличийD)Pub/sub работает только внутри одного процесса, а очередь — между разными серверами
показать ответ и разбор
+B)Pub/sub доставляет событие ВСЕМ подписчикам; очередь — одному из потребителей// разбор: Очередь (point-to-point): сообщение забирает ОДИН из конкурирующих потребителей — так РАСПРЕДЕЛЯЮТ работу (балансировка), каждое сообщение обрабатывается однократно. Pub/sub (публикация-подписка): событие доставляется ВСЕМ подписчикам (fan-out), у каждого своя копия — так РАЗНЫЕ сервисы реагируют на одно событие независимо (заказ создан → уведомить, списать склад, начислить бонусы). В Kafka это выражается через consumer groups: внутри группы — как очередь (баланс по партициям), между группами — как pub/sub (каждая группа читает всё). RabbitMQ — через exchange-типы. Выбор модели — по тому, «поделить работу» или «оповестить многих».
- Зачем сервису публиковать доменное событие (например, OrderCreated)?A)Чтобы сохранить заказ в базе данных: публикация события — это и есть запись в таблицу заказовB)Чтобы синхронно дождаться, пока все заинтересованные сервисы обработают заказ, и только потом ответитьC)Оповестить другие сервисы о факте, не зная и не вызывая их напрямую (слабая связанность)D)Чтобы заменить REST: после перехода на события HTTP-эндпоинты сервису больше не нужны
показать ответ и разбор
+C)Оповестить другие сервисы о факте, не зная и не вызывая их напрямую (слабая связанность)// разбор: Доменное событие сообщает ФАКТ («заказ создан», «платёж прошёл») в прошедшем времени. Издатель просто ПУБЛИКУЕТ его в брокер, НЕ зная и не вызывая конкретных получателей; заинтересованные сервисы ПОДПИСЫВАЮТСЯ и реагируют по-своему (склад резервирует, нотификации шлют письмо, аналитика считает) — это event-driven архитектура со слабой связанностью: добавить нового потребителя = подписаться, не трогая издателя. Плюс события — основа Saga (хореография) и аудита/event sourcing. Важно: публикацию события и запись в свою БД делают согласованно (паттерн transactional outbox), иначе событие может «потеряться» или уйти без коммита.
дальше
Теорию прочитали. Навык ставится повторением
В Сеньорчике эта подтема идёт в ежедневных сессиях: движок возвращает её, пока ответы не станут уверенными, и ведёт прогресс отдельно по каждой подтеме. Теория внутри тоже бесплатна, лимит только на количество вопросов в день.