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.
  • Эксклюзивные очереди со случайным именем — только для отладки и наблюдателей.
  • События приходят повторно и не по порядку. Обработчики делайте идемпотентными и самодостаточными.
Проверьте себя
1. Три сервиса (почта, склад, аналитика) должны получать каждое событие order.created. Как настроить RabbitMQ?
AОдна общая очередь, к которой подключены все три сервиса
BFanout-обменник и три отдельные очереди, каждая привязана к нему
CПубликовать событие три раза в default exchange с разными routing_key
DОдна очередь с prefetch_count=3
2. Продюсер опубликовал событие в fanout-обменник, к которому в этот момент не привязана ни одна очередь. Что случится с сообщением?
AОно подождёт в обменнике, пока не появится подписчик
BБрокер вернёт продюсеру ошибку публикации
CОно исчезнет бесследно — обменник ничего не хранит
DОно попадёт в очередь по умолчанию с именем обменника