Publish/Subscribe: рассылка событий
Одно событие — много независимых реакций. Разбираем fanout: рассылку сообщения всем подписчикам сразу.
Publish/Subscribe — паттерн, при котором продюсер публикует событие, не зная, кто и сколько его прочитает. Копию получает каждый подписчик, а не один из них, как в очереди задач.
Зачем это нужно
Пользователь оформил заказ. Что должно произойти дальше? Отправить письмо с подтверждением. Зарезервировать товар на складе. Записать событие в аналитику. Начислить бонусы. Через месяц продакт попросит ещё и в Telegram-бота уведомление слать.
Наивный вариант — вызывать все эти сервисы из кода оформления заказа по очереди. Через полгода функция create_order знает про пять чужих сервисов, тормозит на самом медленном из них, а падение почтовика роняет оформление заказа целиком. Каждая новая реакция — правка и релиз сервиса заказов.
С fanout сервис заказов делает ровно одно: публикует событие «заказ 7712 создан» в обменник и забывает о нём. Кто на этот обменник подписан — его не касается. Новый потребитель подключается сам, без единой строчки изменений в продюсере. Это то самое «слабое связывание», ради которого люди и ставят брокер.
Продюсер: публикуем событие
import json
import pika
connection = pika.BlockingConnection(pika.ConnectionParameters(host="localhost"))
channel = connection.channel()
# fanout: копию получит каждая привязанная очередь
channel.exchange_declare(exchange="orders.events", exchange_type="fanout", durable=True)
event = {
"type": "order.created",
"order_id": 7712,
"user_id": 91,
"total": 4390,
"items": [{"sku": "TEA-001", "qty": 2}],
}
channel.basic_publish(
exchange="orders.events",
routing_key="", # fanout игнорирует routing_key
body=json.dumps(event).encode(),
properties=pika.BasicProperties(
delivery_mode=pika.DeliveryMode.Persistent,
content_type="application/json",
message_id="order.created:7712", # пригодится для дедупликации
),
)
print("событие order.created опубликовано")
connection.close()
Обратите внимание: продюсер не объявляет ни одной очереди. Он вообще не знает, существуют ли подписчики. Его зона ответственности заканчивается на обменнике — и это принципиально. routing_key для fanout не используется: обменник этого типа копирует сообщение во все привязанные к нему очереди, не глядя на ключ.
Подписчики: у каждого своя очередь
import json
import pika
QUEUE = "orders.email" # у сервиса аналитики будет orders.analytics
connection = pika.BlockingConnection(pika.ConnectionParameters(host="localhost"))
channel = connection.channel()
# оба объявления идемпотентны: подписчик поднимет что нужно, если продюсер ещё не стартовал
channel.exchange_declare(exchange="orders.events", exchange_type="fanout", durable=True)
channel.queue_declare(queue=QUEUE, durable=True)
channel.queue_bind(exchange="orders.events", queue=QUEUE)
def on_event(ch, method, properties, body):
event = json.loads(body)
if event["type"] != "order.created":
ch.basic_ack(delivery_tag=method.delivery_tag)
return
send_email(event["user_id"], event["order_id"])
ch.basic_ack(delivery_tag=method.delivery_tag)
print("письмо по заказу", event["order_id"], "отправлено")
channel.basic_qos(prefetch_count=10)
channel.basic_consume(queue=QUEUE, on_message_callback=on_event)
channel.start_consuming()
Ключевая строка здесь — queue_bind. Она создаёт привязку: связь между обменником и очередью. Каждый подписчик заводит свою собственную очередь и привязывает её к общему обменнику. Сообщение, попавшее в orders.events, брокер размножит по всем привязкам.
Как это работает
Обменник ничего не хранит
Это первое, что ломает интуицию новичка. Exchange — не буфер и не «тема» с историей. Он живёт долю миллисекунды: получил сообщение, посмотрел на список привязок, разложил копии по очередям, забыл. Хранят только очереди.
Отсюда прямое следствие: если к обменнику не привязана ни одна очередь, сообщение просто исчезает. Без ошибки, без предупреждения. Опубликовали событие до того, как подписчик успел стартовать и создать очередь? Событие потеряно. Именно поэтому подписчик, а не продюсер, объявляет и очередь, и привязку, — и делает это при каждом старте.
Если тихая потеря недопустима, публикуйте с флагом mandatory=True и вешайте Basic.Return-колбэк: тогда брокер вернёт продюсеру недоставленное сообщение. Более надёжный вариант — настроить обменнику alternate exchange, куда падает всё, что никуда не легло.
Одна очередь на подписчика — не на всех
Самое опасное заблуждение. Сравните две конфигурации:
| Конфигурация | Что получится |
Три сервиса — три очереди, каждая привязана к orders.events | Каждый сервис получает все события. Это pub/sub. |
| Три сервиса читают одну общую очередь | Событие достанется одному случайному из них. Это work queue, и письмо уйдёт вместо резервирования склада. |
Три инстанса одного сервиса на одной очереди orders.email | Правильно: письмо отправит один инстанс, это масштабирование подписчика. |
Правило простое: очередь — это подписчик-логический-сервис, а не подписчик-процесс. Сколько сервисов слушает событие — столько очередей. Сколько копий сервиса запущено — столько консьюмеров на одной его очереди.
Временные очереди для наблюдателей
Иногда подписчик нужен на пять минут: посмотреть глазами, что реально летит в шину. Для этого есть эксклюзивные очереди со случайным именем — они удаляются, как только соединение закрылось.
result = channel.queue_declare(queue="", exclusive=True) # брокер придумает имя сам
tmp_queue = result.method.queue # например, amq.gen-JzTY20BRgKO...
channel.queue_bind(exchange="orders.events", queue=tmp_queue)
channel.basic_consume(queue=tmp_queue, on_message_callback=print_event, auto_ack=True)
channel.start_consuming()
Для отладки — отлично. Для постоянного сервиса — категорически нет: пока сервис перезапускается, его временная очередь не существует, и все события этого промежутка теряются навсегда. Постоянному подписчику нужна durable-очередь с фиксированным именем: она копит события, пока сервис лежит.
Как убедиться, что подписка живая
Половина «загадочных» багов pub/sub — это отсутствующая привязка. Проверяется за десять секунд из консоли:
rabbitmqctl list_bindings source_name destination_name
rabbitmqctl list_queues name messages consumers
Первая команда покажет, какие очереди реально привязаны к orders.events. Если вашей там нет — сообщения до неё и не долетали. Вторая выдаёт глубину очереди и число консьюмеров. Комбинация «messages растёт, consumers = 0» означает, что подписчик отвалился, а события копятся. Комбинация «messages = 0, consumers = 1» на живой системе — что подписчик успевает, всё в порядке. То же самое видно в веб-интерфейсе на порту 15672, во вкладке Exchanges: там у обменника перечислены все привязки.
Fanout — не единственный вариант
Fanout рассылает всё и всем. Часто этого достаточно: подписчик сам отбрасывает ненужные типы событий, как в примере выше с проверкой event["type"]. Но когда событий десятки типов, а сервис интересуется одним, гонять через него весь поток расточительно — фильтровать должен брокер. Для этого есть direct (точное совпадение ключа) и topic (шаблоны вида order.*.created), о которых речь в разделе про маршрутизацию. Логика подписки при этом не меняется: своя очередь, queue_bind, свой консьюмер — меняется только тип обменника и ключ привязки.
Частые ошибки
- Публикация в default exchange вместо своего. Классика:
basic_publish(exchange="", routing_key="orders.email"). Формально работает, письмо уйдёт — но вы отправили сообщение конкретной очереди, и никакого pub/sub нет. Второй подписчик не появится никогда. - Подписчик не объявляет привязку. Понадеялись, что «очередь уже создали руками в UI». После переезда на новый инстанс брокера всё молча перестаёт работать. Объявляйте exchange, queue и bind в коде при старте — эти операции идемпотентны.
- Ожидание порядка между подписчиками. Аналитика может увидеть
order.createdна секунду позже склада, а если у подписчика несколько консьюмеров — события даже внутри одного сервиса могут обработаться не по порядку. Не стройте логику на «сначала придёт A, потом B»; кладите в событие всё нужное для самостоятельной обработки. - Забыли про идемпотентность. При переподключении воркера сообщение придёт повторно, и клиент получит два одинаковых письма. Сохраняйте
message_idобработанных событий и отбрасывайте дубли. - Мёртвый подписчик копит очередь. Сервис аналитики выключили полгода назад, а его durable-очередь всё ещё привязана к обменнику и растёт. Брокер упрётся в диск. Отвязывайте (
queue_unbind) и удаляйте очереди отключённых подписчиков — либо ставьте им TTL и лимит длины. - Толстые события. Fanout размножает тело сообщения на каждую очередь: событие на 2 МБ и десять подписчиков — это 20 МБ в памяти брокера. Публикуйте идентификаторы и минимум полей.
Итоги
- Fanout-обменник копирует сообщение во все привязанные очереди;
routing_keyне влияет ни на что. - Продюсер знает только обменник. Подписчики появляются и исчезают без его ведома.
- Один логический подписчик = одна durable-очередь +
queue_bind. Общая очередь на разные сервисы превращает pub/sub в work queue. - Обменник ничего не хранит: без привязанных очередей событие исчезает бесследно. Спасают
mandatory=Trueи alternate exchange. - Эксклюзивные очереди со случайным именем — только для отладки и наблюдателей.
- События приходят повторно и не по порядку. Обработчики делайте идемпотентными и самодостаточными.