Publish/subscribe
Pub/sub меняет доставку один-к-одному на веерную рассылку по топикам: издатель публикует в топик, каждый заинтересованный consumer получает копию. Глубже всего раскол между брокерами-очередями, что удаляют при потреблении, и log-системами с append-only логом и реплеем.
Когда оформлялся заказ, сервис заказов по очереди дёргал четыре других: склад, 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 и хочет ли реплей.
- 01Что pub/sub добавляет сверх очереди точка-точка и когда что использовать?
- 02Сравни брокерный pub/sub с log-based и что следует из выбора.
- 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-уровень. Открой, попробуй, потом открой ответ.
Что-то непонятно?
Задай вопрос по этому уроку. Вопросы анонимны и попадают напрямую автору — урок станет лучше.