open atlas
↑ К треку
Основы System Design SD · 06 · 02

Publish/subscribe

Pub/sub меняет доставку один-к-одному на веерную рассылку по топикам: издатель публикует в топик, каждый заинтересованный consumer получает копию. Глубже всего раскол между брокерами-очередями, что удаляют при потреблении, и log-системами с append-only логом и реплеем.

SD Middle ◷ 18 min
Уровень
ОсновыJuniorMiddleSenior

Когда оформлялся заказ, сервис заказов по очереди дёргал четыре других: склад, email, аналитику, антифрод. Каждая новая фича, которой был интересен заказ, означала ещё один вызов, вшитый в сервис заказов, ещё один таймаут для обработки, ещё один деплой, рискующий сломать checkout. Сервис заказов стал хабом, знающим про всех. Команда перевернула это: сервис заказов теперь публикует одно событие «заказ оформлен» в топик и забывает, кто слушает. Склад, email, аналитика и антифрод подписываются независимо. Добавить пятого consumer’а — сервис лояльности — теперь это деплой только сервиса лояльности; сервис заказов не меняется. Одна запись, много независимых читателей, ноль связности между ними.

От один-к-одному к один-ко-многим

Очередь точка-точка доставляет каждое сообщение одному consumer’у. Publish/subscribe доставляет каждое сообщение каждой заинтересованной стороне. Издатель публикует сообщение в именованный топик; брокер копирует его в каждую подписку; каждый подписчик потребляет свою копию независимо. Издатель называет топик, а не получателя — он понятия не имеет, кто слушает, и consumer’ы приходят и уходят, а издатель этого не знает и об этом не заботится. Эта косвенность и есть весь смысл: она добавляет уровень развязки сверх временной развязки обычной очереди.

Результат — инверсия из хука. Вместо producer’а, держащего список из N downstream’ов для вызова (и деплоя при каждом изменении списка), producer публикует один раз, а consumer’ы решают, что им интересно. Это разница между командами («отправь это письмо» — адресовано одному обработчику, задача → используй очередь) и событиями («заказ оформлен» — факт, разосланный всем, кому интересно → используй pub/sub). Правильно провести это различие — это бОльшая часть дизайна: у команды ровно один корректный обработчик; у события ноль или много, и нельзя зашивать consumer’ов в producer.

Две архитектуры под одним именем: брокер против лога

«Pub/sub» описывает поведение (fan-out), и его реализуют два очень разных механизма. Этот раскол — самое важное, что надо понять в этом уроке.

Брокерный стиль (RabbitMQ, ActiveMQ, AWS SNS→SQS). Брокер маршрутизирует копию каждого сообщения в пер-подписчиковую очередь, и как только подписчик подтвердил, эта копия удаляется. Хранение сообщений транзиентно — задача брокера в доставке, а не в удержании. Подписчик, что не был подключён в момент публикации, может попросту никогда не увидеть сообщение (если ты не настроил durable-подписки). Брокер отслеживает состояние доставки за consumer’ов.

Log-based (Apache Kafka, AWS Kinesis, Apache Pulsar). Топик — это append-only лог: сообщения пишутся по порядку и удерживаются на настроенное окно (часы, дни или навсегда), независимо от того, прочитал ли их кто-то. Сообщения не пушатся consumer’ам с удалением; они тянут (pull) и отслеживают собственную позицию — офсет (offset) — в логе. Несколько независимых consumer’ов читают один лог на своих офсетах, не мешая друг другу. Чтение не потребляет; лог просто помнит, где каждый consumer.

БРОКЕР (RabbitMQ/SNS)            ЛОГ (Kafka/Kinesis)
publish → веер копий → очереди   append → неизменяемый лог [m0 m1 m2 m3 m4 …]
consumer ack → копия УДАЛЕНА      consumer A читает, двигает свой ОФСЕТ → 3
состояние живёт в БРОКЕРЕ         consumer B читает независимо → офсет 1
«труба доставки»                 состояние (офсет) живёт в CONSUMER'е
                                  «долговечная, реплеябельная запись»

Этот один архитектурный выбор каскадом задаёт всё остальное: реплей, порядок, пропускную способность и операционную модель — всё следует из «удалить при потреблении» против «append-only лог с офсетами».

Почему это работает

Почему log-модель делает реплей почти бесплатным, а брокерная не может? В логе сообщение не удаляется при чтении — двигается только офсет consumer’а. Поэтому «переиграть последний час» — это просто «сбрось мой офсет на час назад и читай вперёд снова»; данные всё ещё там. Совершенно новый consumer может стартовать с офсета 0 и переобработать всю историю. В брокере сообщение исчезает в момент подтверждения: реиграть нечего, потому что удержание никогда не было задачей брокера. Именно поэтому event sourcing, перестройка downstream’а с нуля и добавление consumer’а, которому нужен бэкфилл, — всё это естественно на логе и неловко-до-невозможно на классическом брокере. И поэтому «нужна ли мне возможность переобработать историю?» — вопрос, решающий, к чему ты тянешься.

Consumer-группы и порядок по партициям

Лог-топик разбит на партиции для параллелизма — единицу и масштаба, и порядка. Consumer-группа — это набор кооперирующих consumer’ов, делящих работу по топику: каждая партиция назначена ровно одному consumer’у в этой группе, поэтому добавление consumer’ов вплоть до числа партиций масштабирует пропускную способность линейно. Важно: разные группы каждая получают полный поток — группа «аналитика» и группа «антифрод» обе читают каждое сообщение на своих офсетах. Так один топик кормит много независимых подписчиков (fan-out между группами), при этом каждый подписчик балансирует нагрузку внутри себя (деление внутри группы).

Порядок следует из партиций. Лог гарантирует порядок внутри партиции, а не по всему топику. Сообщения с одним ключом (например, user_id) хешируются в одну партицию, поэтому все события одного пользователя остаются упорядочены относительно друг друга; события разных пользователей могут чередоваться. Это тот же компромисс порядка по ключу из урока про очереди, теперь сделанный структурным: ты получаешь параллелизм между партициями и порядок внутри одной — и выбираешь ключ партиционирования, чтобы класть «то, что должно остаться упорядоченным» вместе. Выберешь ключ со слишком малым числом различных значений — создашь горячую партицию, которую нельзя распараллелить.

Так лог — это просто навороченная очередь? Нет.

Самый глубокий вывод: очередь и лог — не один инструмент с разным брендингом. Очередь — это труба доставки: сообщения протекают и исчезают после обработки; её состояние — «что осталось сделать». Лог — это долговечная запись фактов: сообщения остаются, а consumer’ы движутся по ним; его состояние — «где каждый читатель». Очередь отвечает на «какая работа осталась?»; лог отвечает на «что произошло и могу ли я прочитать это снова?». К брокеру/очереди тянешься, когда работу надо сделать раз и забыть, а к логу — когда события суть факты, которые ты можешь захотеть разослать веером, переставить consumer’ов над ними или переиграть. Путать их — трактовать Kafka как очередь задач, из которой удаляешь, или очередь как event store, который реплеишь, — частая и дорогая несостыковка.

Выбери лучший вариант

Новый сервис обнаружения мошенничества должен перепотребить все события заказов за последние 30 дней для бэкфилла модели, затем продолжить получать живые события. Сервис лояльности тоже независимо нуждается в том же живом потоке. Какая архитектура мессаджинга подходит для обоих требований?

Частая ошибка

Частая ловушка fan-out: считать, что «pub/sub» означает, что каждый consumer всегда видит каждое сообщение независимо от того, когда присоединился. На брокере подписчик, что офлайн в момент публикации, полностью пропускает сообщение, если ты явно не настроил durable-подписку — поэтому новый сервис аналитики получает ноль истории, только сообщения, опубликованные после подписки, и тихо недосчитывается. На логе риск обратный: consumer, что лежал дольше окна удержания, возвращается и обнаруживает, что его офсет упал за конец удержанных данных и надо перепрыгнуть вперёд, теряя сообщения. Оба — баги «я думал, получу всё»: всегда спрашивай, что происходит с consumer’ом, который отсутствует, медленен или совсем новый.

Викторина

Нужно добавить совершенно новый consumer аналитики, который должен переобработать события заказов за последние 30 дней, чтобы заполнить дашборд. Какая архитектура поддерживает это напрямую и почему?

Викторина

На топике Kafka с 4 партициями нужно обрабатывать события каждого пользователя строго по порядку, но распараллеливать между пользователями. Как ключевать и потреблять?

Закончи аналогию

Глубочайшая разница между очередью и логом: очередь — труба доставки, чьи сообщения исчезают после обработки, тогда как лог удерживает сообщения, и каждый consumer отслеживает свой собственный _______ в нём — ровно поэтому реплей и добавление новых consumer'ов first-class на логе.

Это высота композиции: выбор fan-out и выбор лог-против-брокера для дизайна. Отдельный трек queues прорабатывает внутренности — лидерство партиций Kafka и ISR, типы exchange и биндинги RabbitMQ, точные протоколы коммита офсета. Используй его, когда эксплуатируешь конкретный брокер; оставайся здесь, когда решаешь, хочет ли дизайн fan-out и хочет ли реплей.

Вспомните перед уходом
  1. 01
    Что pub/sub добавляет сверх очереди точка-точка и когда что использовать?
  2. 02
    Сравни брокерный pub/sub с log-based и что следует из выбора.
  3. 03
    Как consumer-группы и партиции дают и fan-out, и порядок?
Итог

Publish/subscribe превращает доставку один-к-одному в веерную рассылку (fan-out) по топикам: издатель публикует в именованный топик, брокер копирует его в каждую подписку, и издатель называет топик, а не получателя — так producer и consumer’ы развязаны, и consumer добавляется без правки producer’а. Используй очередь для команд (один обработчик, сделать-раз) и pub/sub для событий (факты, разосланные нулю-или-многим). Хребет урока в том, что имя носят две архитектуры: брокерный стиль (RabbitMQ/SNS) маршрутизирует копии и удаляет при ack — транзиентно, труба доставки — тогда как log-based (Kafka/Kinesis) держит append-only лог, удержанный на окно, где каждый consumer отслеживает свой офсет — долговечная, реплеябельная запись. Этот один выбор каскадирует в реплей (почти бесплатен на логе, невозможен на брокере), онбординг новых consumer’ов и эксплуатацию. Consumer-группы дают fan-out между группами и деление нагрузки внутри одной (партиция на consumer’а, лишние простаивают), а партиции дают порядок внутри партиции (ключуй по сущности) плюс параллелизм между ними. Вывод, который стоит держать: лог — не навороченная очередь — очередь отвечает «какая работа осталась», лог отвечает «что произошло и могу ли прочитать снова». Отдельный трек queues покрывает внутренности брокеров; оставайся здесь, чтобы решить fan-out и реплей. Теперь, когда встретишь новый сервис аналитики, которому нужна история, — первым делом спроси: это лог с достаточным удержанием или брокер, что уже удалил всё нужное?

Практика

Начни сверху. Задачи идут от простого к сложному: вспомнить факт, применить к случаю, затем senior-уровень. Открой, попробуй, потом открой ответ.

вспомнитьприменитьуглубить0 из 8 завершено
Связанные уроки

Что-то непонятно?

Задай вопрос по этому уроку. Вопросы анонимны и попадают напрямую автору — урок станет лучше.

хоткеи развернуть
поиск
K
пред. пьеса
k
след. пьеса
j
тиры
t
это меню
?
sources3
expand
  1. 01
  2. 02
  3. 03

Trademarks belong to their respective owners. Editorial reference only.