Перейти к основному контенту
Tech Path Finder
КурсыИнтервьюКод-ревьюБлог
Tech Path Finder

Персонализированный путеводитель в IT. Квизы, мок-интервью, код ревью и аналитика прогресса.

@potapov_me

Платформа

  • Курсы
  • Прогресс
  • Мок-интервью
  • Код ревью
  • Живое ревью с ИИ
  • Тренажёр переговоров
  • Закладки

Контент

  • Блог
  • Главная
  • Обратная связь

Компания

  • О проекте
  • Тарифы
  • Условия использования
  • Конфиденциальность
  • Согласие на обработку данных
  • Cookie
  • Реквизиты

Аккаунт

  • Войти
  • Зарегистрироваться
  • Профиль

© 2026 Tech Path Finder. Все права защищены.

·ИП Потапов К.С.·Политика конфиденциальности·
Сделано с ❤️ в России
  1. Публикация сообщений (Publishing)
publishing

Публикация сообщений (Publishing)

Отправка сообщений во все брокеры: базовая публикация, кастомизация заголовков, delivery modes, подтверждение отправки.

Публикация сообщений (Publishing)

Отправка сообщений во все брокеры: базовая публикация, кастомизация заголовков, delivery modes, подтверждение отправки.

#Архитектура публикации в FastStream

FastStream абстрагирует различия между брокерами сообщений, предоставляя единый API для публикации. Это позволяет:

  • Сменить брокер без изменения бизнес-логики — код публикации остается тем же
  • Использовать специфичные фичи брокера — когда нужны уникальные возможности
  • Генерировать AsyncAPI документацию — автоматически из кода

Почему это важно: В production вы можете начать с RabbitMQ для простоты, но позже мигрировать на Kafka для throughput. FastStream делает это безболезненным.

#Базовая публикация

#Простая публикация

from faststream import FastStream from faststream.rabbit import RabbitBroker broker = RabbitBroker("amqp://localhost") app = FastStream(broker) # Публикация словаря await broker.publish({"order_id": 123}, "orders") # Публикация строки await broker.publish("Hello, World!", "greetings") # Публикация Pydantic-модели from pydantic import BaseModel class Order(BaseModel): id: int amount: float order = Order(id=123, amount=99.99) await broker.publish(order, "orders") # Автоматическая сериализация в JSON

Почему Pydantic модели предпочтительнее словарей:

  • Валидация на отправке — ошибка обнаружится до публикации, а не в обработчике
  • Автодокументирование — AsyncAPI сгенерирует схему сообщения
  • Type safety — IDE подскажет поля, меньше багов

Когда использовать строки: Только для интеграции с legacy системами, ожидающими raw текст.

#Публикация из обработчика

@broker.subscriber("orders.new") async def process_order(order: dict): # Обработать заказ await db.save(order) # Опубликовать событие await broker.publish( {"order_id": order["id"], "status": "processed"}, "orders.processed" )

Почему это анти-pattern: Если db.save() выполнится, а publish() упадёт — данные в БД есть, а событие о том не сообщает. Подписчики orders.processed не узнают о заказе.

Правильный подход — транзакционная публикация:

@broker.subscriber("orders.new") async def process_order(order: dict): # Сначала готовим событие event = {"order_id": order["id"], "status": "processed"} # Сохраняем в БД вместе с событием (outbox pattern) await db.save_with_event(order, event) # Публикуем после успешного коммита await broker.publish(event, "orders.processed")

Или используйте транзакционные outbox паттерн для гарантии доставки.

#Publisher декоратор

Декоратор @publisher объявляет издатель на уровне приложения:

@broker.publisher("orders.processed") @broker.subscriber("orders.new") async def process_order(order: dict): await db.save(order) # publisher автоматически доступен await process_order.publish( {"order_id": order["id"], "status": "processed"} )

Преимущества перед broker.publish():

  • Декларативность — все publishers видны в сигнатуре функции
  • Автогенерация AsyncAPI — документация включает все channels
  • Тестируемость — легко мокать process_order.publish в unit тестах
  • Response routing — можно возвращать значение, и оно автоматически публикуется

Когда использовать broker.publish() напрямую:

  • Динамические routing key (зависят от данных)
  • Публикация из cron jobs, startup hooks
  • Batch публикация (несколько сообщений за раз)

Пример response routing:

@broker.subscriber("orders.new") @broker.publisher("orders.processed") async def process_order(order: dict) -> dict: # Возвращаемое значение автоматически публикуется return {"order_id": order["id"], "status": "processed"}

#Кастомизация сообщений

#Заголовки (Headers)

await broker.publish( {"order_id": 123}, "orders", headers={ "priority": "high", "trace_id": "abc-123", "content-type": "application/json" } )

Зачем нужны заголовки:

Заголовки — это метаданные, которые не являются частью бизнес-логики сообщения. Они решают несколько критичных задач:

  1. Distributed tracing — trace_id и span_id позволяют отследить путь запроса через все микросервисы
  2. Routing — брокеры могут роутить сообщения на основе заголовков (например, RabbitMQ headers exchange)
  3. Версионирование — schema_version: "1.0" помогает обработчикам поддерживать обратную совместимость
  4. Multi-tenancy — tenant_id для изоляции данных разных клиентов в одной очереди
  5. Observability — мониторинговые системы собирают метрики по заголовкам

Production пример с tracing:

import uuid from contextvars import ContextVar # Context variable для хранения trace_id в рамках обработки trace_id_ctx: ContextVar[str] = ContextVar("trace_id", default="") @broker.subscriber("orders.new") async def process_order(order: dict, message: Message): # Получаем trace_id из входящего сообщения или генерируем новый trace_id = message.headers.get("trace_id", str(uuid.uuid4())) trace_id_ctx.set(trace_id) # Логируем с trace_id — теперь все логи этого запроса связаны logger.info("Processing order", extra={"trace_id": trace_id}) # Публикуем с тем же trace_id await broker.publish( {"order_id": order["id"], "status": "processed"}, "orders.processed", headers={"trace_id": trace_id} # Сохраняем trace chain )

Когда не использовать заголовки:

  • Бизнес-данные (payload сообщения) — это belongs в теле
  • Большие бинарные данные (заголовки имеют лимиты размера)
  • Чувствительные данные без шифрования (заголовки могут логироваться)

#Correlation ID

Для трассировки и RPC:

import uuid correlation_id = str(uuid.uuid4()) await broker.publish( {"request": "data"}, "requests", correlation_id=correlation_id ) # В обработчике @broker.subscriber("requests") async def handle_request(data: dict, message: Message): trace_id = message.correlation_id # Логировать с trace_id logger.info(f"Processing {data}", extra={"trace_id": trace_id})

Correlation ID vs Trace ID — в чем разница:

  • Trace ID — отслеживает весь путь запроса через ВСЕ микросервисы (end-to-end)
  • Correlation ID — связывает ЗАПРОС-ОТВЕТ между двумя сервисами (request-response pair)

Когда использовать correlation_id:

  1. RPC паттерн — сервис A запрашивает сервис B, нужен ID для связки запрос-ответ
  2. Дедупликация — если сообщение доставлено дважды (at-least-once delivery), correlation_id помогает отфильтровать дубли
  3. Audit trail — compliance требует отследить, почему было принято решение (все связанные сообщения)

Production пример с дедупликацией:

import uuid from redis import AsyncRedis dedup_cache = AsyncRedis("redis://localhost") @broker.subscriber("payments.process") async def process_payment(payment: dict, message: Message): corr_id = message.correlation_id # Проверяем, не обрабатывали ли уже это сообщение if await dedup_cache.exists(f"processed:{corr_id}"): logger.warning(f"Duplicate message {corr_id}, skipping") return # Обрабатываем платеж await process(payment) # Запоминаем correlation_id (TTL 24 часа) await dedup_cache.set(f"processed:{corr_id}", "1", ex=86400) # Публикуем результат await broker.publish( {"payment_id": payment["id"], "status": "done"}, "payments.done", correlation_id=corr_id # Сохраняем correlation )

#Reply To

Для RPC указать очередь ответа:

await broker.publish( {"query": "select *"}, "queries", reply_to="query-results" )

Почему RPC через брокер — это не always good idea:

RPC (request-response) через брокер сообщений увеличивает latency по сравнению с HTTP/gRPC. Используйте когда:

  • Нужна асинхронность — клиент не ждет блокируя
  • Есть service mesh ограничения — HTTP недоступен
  • Требуется гарантированная доставка — брокер retry'ит при ошибках

Для low-latency сценариев используйте HTTP/gRPC напрямую.

#Delivery modes

#Персистентность (RabbitMQ)

# Персистентное сообщение (сохраняется на диск) await broker.publish( {"order_id": 123}, "orders", persistent=True # delivery_mode=2 ) # Неперсистентное (только в памяти) await broker.publish( {"event": "click"}, "events", persistent=False # delivery_mode=1 )

Что происходит под капотом:

  • persistent=True — RabbitMQ записывает сообщение на диск перед подтверждением отправки. При перезапуске брокера сообщение сохраняется.
  • persistent=False — сообщение только в RAM. Быстрее, но при креше брокера теряется.

Критичный нюанс: Персистентность сообщения работает ТОЛЬКО если очередь и exchange тоже персистентные. Иначе сообщение сохраняется, но queue/exchange исчезают — сообщение теряется.

Trade-off: Throughput vs Durability

ModeThroughputDurabilityUse Case
persistent=False~50K msg/sПри креше теряетсяМетрики, клики, логи
persistent=True~5-10K msg/sПереживает крешЗаказы, платежи, регистрации

Production правило:

  • persistent=True: критичные данные (заказы, платежи, пользовательские данные)
  • persistent=False: временные события (клики, метрики, health checks), которые можно regenerate

#Priority (Приоритет)

# Высокий приоритет await broker.publish( {"order_id": 123, "vip": True}, "orders", priority=10 # 0-10, чем выше — тем приоритетнее ) # Низкий приоритет await broker.publish( {"order_id": 124}, "orders", priority=1 )

Как это работает:

RabbitMQ хранит сообщения в queue. При priority брокер сортирует очередь — сообщения с высоким приоритетом доставляются ПЕРВЫМИ, даже если они пришли позже.

Критичное требование: Очередь ДОЛЖНА быть объявлена с max_priority, иначе приоритет игнорируется:

# Без max_priority приоритет НЕ работает broker = RabbitBroker("amqp://localhost") @broker.subscriber( "orders", queue_options={ "x-max-priority": 10 # Обязательно! } ) async def process_order(order: dict): ...

Когда использовать приоритеты:

  • VIP клиенты — платежи премиум клиентов обрабатываются первыми
  • SLA таймауты — сообщения с истекающим TTL получают высокий приоритет
  • Retry с priority — первые попытки низкий приоритет, последние (перед alert) высокий

Когда НЕ использовать:

  • Высокий throughput — priority queue ~20% медленнее обычной
  • Fair scheduling — если все клиенты равны, приоритеты создают starvation
  • Kafka — Kafka НЕ поддерживает приоритеты (используйте partition key)

Production паттерн: Dynamic priority based on SLA

from datetime import datetime, timedelta @broker.subscriber("orders.new") async def process_order_with_priority(order: dict): created_at = datetime.fromisoformat(order["created_at"]) age_minutes = (datetime.utcnow() - created_at).total_seconds() / 60 # Чем старше сообщение, тем выше приоритет if age_minutes > 50: # Почти timeout (60 мин) priority = 10 elif age_minutes > 30: priority = 5 else: priority = 1 await broker.publish( order, "orders.processing", priority=priority )

#TTL (Time To Live)

# Сообщение истекает через 1 час await broker.publish( {"cache_key": "user:123"}, "cache-invalidation", expiration=3600000 # миллисекунды )

Зачем нужен TTL:

TTL гарантируетет, что сообщение будет удалено из очереди после указанного времени, даже если не было обработано.

Use cases:

  1. Cache invalidation — устаревшие данные уже не актуальны, нет смысла обрабатывать
  2. Real-time notifications — push уведомление через 2 часа после события уже не нужно
  3. Circuit breaker — если очередь забита, старые сообщения expire'ятся, давая системе восстановиться
  4. Rate limiting — сообщения с коротким TTL не создают backlog

TTL на уровне queue vs message:

# TTL на очереди (все сообщения живут 1 час) @broker.subscriber( "notifications", queue_options={ "x-message-ttl": 3600000 # 1 час для всех } ) async def handle_notification(msg: dict): ... # TTL на сообщение (гибче) await broker.publish( {"type": "urgent", "data": "..."}, "notifications", expiration=60000 # 1 минута для критичных ) await broker.publish( {"type": "normal", "data": "..."}, "notifications", expiration=3600000 # 1 час для обычных )

Production нюанс: Сообщение с истекшим TTL не удаляется сразу. Оно удаляется когда:

  • Достигает head очереди (при delivery consumer'у)
  • При purge очереди
  • При restart брокера

Это значит, что при большом backlog вы можете получить сообщения, которые "expired" но еще не удалены.

#Публикация в разные брокеры

#Как выбрать брокер для задачи

ХарактеристикаRabbitMQKafkaRedisNATS
Гарантия доставкиAt-least-onceExactly-once (с idempotence)At-most-once (pub/sub)At-most-once (core)
Throughput~10K msg/s~1M msg/s~50K msg/s~100K msg/s
Latency<10ms<5ms<1ms<1ms
Message orderingQueue levelPartition levelNo guaranteeNo guarantee (core)
Message replay❌✅✅ (streams)✅ (JetStream)
Best forTask queues, RPCEvent sourcing, streamsCaching, simple pub/subIoT, real-time

#RabbitMQ

from faststream.rabbit import RabbitBroker broker = RabbitBroker("amqp://localhost") await broker.publish( {"order_id": 123}, "orders.new", # routing key exchange="orders", # exchange name exchange_type="direct", # тип exchange persistent=True, # персистентность priority=5, # приоритет headers={"tenant": "acme"} # заголовки )

Почему RabbitMQ — good default choice:

  • Flexible routing — exchanges (direct, topic, fanout, headers) для сложных routing паттернов
  • Per-message TTL — контроль времени жизни
  • Priority queues — обработка по приоритетам
  • Dead letter exchanges — автоматический retry/failed handling
  • Management UI — мониторинг из коробки

Когда НЕ RabbitMQ:

  • Высокий throughput (>100K msg/s) — выбирайте Kafka
  • Message replay — Kafka хранит историю, RabbitMQ нет
  • Simple pub/sub — Redis проще для базовых задач

#Kafka

from faststream.kafka import KafkaBroker broker = KafkaBroker("localhost:9092") await broker.publish( {"order_id": 123}, "orders", # topic key=b"order-123", # partition key (для порядка) headers={"trace_id": "abc"}, # заголовки partition=0 # конкретная партиция (редко) )

Partition key — это критично:

Partition key определяет, в какую партицию попадет сообщение. Все сообщения с одинаковым key идут в одну партицию, что гарантирует:

  1. Ordering — сообщения для одного агрегата обрабатываются в порядке отправки
  2. Parallelism — разные агрегаты обрабатываются параллельно

Production примеры partition key:

# Все события одного пользователя в одной партиции (ordering guaranteed) await broker.publish( {"event": "user.login"}, "user-events", key=b"user-123" # user_id как partition key ) # Все заказы одного ресторана в одной партиции await broker.publish( {"order_id": 456}, "orders", key=b"restaurant-789" # restaurant_id как partition key )

Когда НЕ указывать key:

  • События не требуют ordering (метрики, логи)
  • Хотите равномерное распределение по партициям
  • Риск: Без key сообщения одного агрегата могут попасть в разные партиции — ordering нарушен!

Kafka vs RabbitMQ trade-offs:

AspectKafkaRabbitMQ
OrderingPer-partitionPer-queue
Replay✅ Хранит историю❌ Нет replay
Throughput~1M msg/s~10K msg/s
RoutingSimple (topics)Flexible (exchanges)
Latency<5ms<10ms

#Redis

from faststream.redis import RedisBroker broker = RedisBroker("redis://localhost:6379") # Pub/Sub await broker.publish( {"event": "user.login"}, "events" ) # Stream (персистентность) await broker.publish( {"event": "user.login"}, "events", stream=True ) # List (очередь задач) await broker.publish( {"task": "send_email"}, "tasks", list=True )

Redis — это НЕ типичный брокер сообщений:

Redis — in-memory datastore, который "из коробки" поддерживает pub/sub, streams и lists. Это одновременно и сила, и слабость.

Три паттерна в Redis:

PatternDurabilityOrderingUse Case
Pub/Sub❌ At-most-once❌Real-time notifications, live updates
Stream✅ Persistent✅ Per-streamEvent logging, simple event sourcing
List✅ Persistent✅ FIFOTask queues, job processing

Когда Redis — правильный выбор:

  • Уже есть Redis в инфраструктуре — не нужно ставить отдельный брокер
  • Simple pub/sub — live notifications, чаты, real-time дашборды
  • Task queue —Celery использует Redis lists для очередей задач
  • Low latency — sub-millisecond задержка

Когда НЕ Redis:

  • Критичные данные — Redis не дает гарантий доставки (pub/sub)
  • Complex routing — нет exchanges, topic routing как в RabbitMQ
  • High availability — Redis Cluster сложен в настройке
  • Message replay — streams поддерживают, но менее гибкие чем Kafka

Production пример: когда использовать stream vs pub/sub

# Pub/Sub — для real-time уведомлений (не критично если потеряется) await broker.publish( {"user_id": 123, "notification": "New message!"}, "notifications", # Pub/Sub channel # Без stream=True — получатели offline = сообщение потеряно ) # Stream — для event logging (нужна персистентность) await broker.publish( {"user_id": 123, "event": "message.sent", "message_id": 456}, "user-events", # Stream name stream=True, # Теперь это Redis Stream — сохраняется )

#NATS

from faststream.nats import NatsBroker broker = NatsBroker("nats://localhost:4222") # Core (без персистентности) await broker.publish( {"event": "user.login"}, "events.user.login" ) # JetStream (персистентность) await broker.publish( {"event": "user.login"}, "events.user.login", stream="events" )

NATS — lightweight alternative:

NATS создан для high-performance IoT и real-time систем. Два режима:

ModeDurabilityOrderingUse Case
Core NATS❌ At-most-once❌IoT telemetry, real-time sync
JetStream✅ Persistent✅ Per-streamEvent sourcing, task queues

Почему NATS:

  • Ultra-low latency — <1ms, быстрее Kafka и RabbitMQ
  • Simple — меньше концепций чем RabbitMQ (нет exchanges, bindings)
  • Auto-TLS — из коробки шифрование
  • Lightweight — бинарник <20MB, минимум RAM
  • IoT protocols — MQTT, WebSockets из коробки

Когда NATS — правильный выбор:

  • IoT devices — тысячи устройств, низкая латентность критична
  • Real-time sync — collaborative editing, live дашборды
  • Service mesh — коммуникация между микросервисами с low latency
  • Простота — нужна минимальная конфигурация

Когда НЕ NATS:

  • Enterprise integration — меньше интеграций чем Kafka/RabbitMQ
  • Message replay — JetStream есть, но менее зрелый чем Kafka
  • Ecosystem — меньше тулов, комьюнити

Production пример: Core vs JetStream

# Core NATS — для real-time telemetry (не критично если потеряется) await broker.publish( {"device_id": "sensor-123", "temperature": 23.5}, "iot.telemetry" # Без stream — device offline = данные потеряны ) # JetStream — для критичных событий await broker.publish( {"device_id": "sensor-123", "alert": "temperature > threshold"}, "iot.alerts", stream="iot-events" # Теперь персистентно )

#Batch публикация

Отправка нескольких сообщений в одном запросе:

#Зачем нужен batch

Network overhead: Каждое сообщение = отдельный network round-trip + сериализация + подтверждение. При 1000 msg/s это создает значительный overhead.

Batch решает: Отправляем N сообщений за один network call → меньше latency, выше throughput.

#RabbitMQ

from faststream.rabbit import RabbitBatch async with RabbitBatch(broker) as batch: batch.publish({"id": 1}, "orders") batch.publish({"id": 2}, "orders") batch.publish({"id": 3}, "orders") # Один network round-trip вместо трёх

Как работает RabbitBatch:

  1. Открывает publisher confirms channel
  2. Буферизует сообщения в памяти
  3. При exit из context manager отправляет все за раз
  4. Ждет подтверждение от брокера

Когда использовать batch:

  • High throughput — >10K msg/s, нужен batch для производительности
  • Bulk operations — импорт данных, миграции
  • Event sourcing — публикация множества events за раз

Batch trade-offs:

AspectSingle PublishBatch Publish
Latency per msg<10ms<10ms (same)
Throughput~10K msg/s~50-100K msg/s
MemoryLowHigher (буферизация)
Failure handlingSingle msg retryВесь batch retry

Production пример: bulk import

async def import_products(products: list[dict]): """Импорт тысяч товаров в систему""" async with RabbitBatch(broker, batch_size=100) as batch: for product in products: batch.publish( product, "products.imported", exchange="products", persistent=True ) # Каждые 100 сообщений отправляются за раз

#Kafka

Kafka по умолчанию batch'ит сообщения на уровне producer'а:

from faststream.kafka import KafkaBroker broker = KafkaBroker( "localhost:9092", # Kafka batch настройки batch_size=16384, # 16KB max batch size linger_ms=5, # Ждать 5ms для заполнения batch ) # Каждое publish вызов может попасть в batch await broker.publish({"id": 1}, "orders") await broker.publish({"id": 2}, "orders") await broker.publish({"id": 3}, "orders") # Kafka автоматически объединит в batch

Kafka batch vs RabbitMQ batch:

AspectKafkaRabbitMQ
Batch automatic✅ Автоматически❌ Явный batch context
Configurationlinger_ms, batch_sizeContext manager
Throughput~1M msg/s~50-100K msg/s

#Подтверждение публикации

#RabbitMQ: подтверждение от брокера

# По умолчанию broker подтверждает публикацию result = await broker.publish( {"order_id": 123}, "orders", raise_timeout=True # Выбросить exception при таймауте ) if result: print("Message published successfully")

#Kafka: подтверждение от лидеров партиций

# acks="all" — подтверждение от всех реплик broker = KafkaBroker( "localhost:9092", acks="all", # Надёжность enable_idempotence=True # Без дубликатов ) await broker.publish({"order_id": 123}, "orders") # Exception при ошибке публикации

Kafka acks — trade-off между reliability и latency:

acks valueЧто подтверждаетLatencyDurabilityUse Case
acks=0Нет подтвержденияLowest❌ Может потерятьМетрики, логи
acks=1Leader толькоLow⚠️ Если leader crashНе критичные данные
acks=allLeader + все ISRHighest✅ Переживает crashЗаказы, платежи

Idempotence producer — почему это важно:

Без idempotence при network timeout producer может отправить сообщение дважды (не уверен, дошло ли первое):

Producer                        Kafka
   |                                |
   |-- Publish message ----------->|
   |                                | (сохраняет)
   |  <-- Network timeout ---------| (ACK потерялся)
   |                                |
   |-- Publish message ----------->| (retry — дубликат!)

С enable_idempotence=True:

  • Producer добавляет sequence number к каждому сообщению
  • Kafka отфильтровует дубликаты по sequence number
  • Exactly-once semantics гарантируется

Когда НЕ включать idempotence:

  • Максимальная производительность — ~5% overhead
  • Kafka < 0.11 — idempotence не поддерживается

#Пример: публикация событий домена

Этот пример показывает production-ready паттерн для event-driven архитектуры с domain events:

from faststream import FastStream from faststream.rabbit import RabbitBroker from pydantic import BaseModel from datetime import datetime import uuid class DomainEvent(BaseModel): id: str type: str aggregate_id: str timestamp: str payload: dict broker = RabbitBroker("amqp://localhost") app = FastStream(broker) # Декларация publishers @broker.publisher("orders.events") @broker.publisher("users.events") @broker.publisher("payments.events") class EventPublisher: pass async def publish_event(event: DomainEvent, exchange: str): """Универсальная функция публикации событий""" await broker.publish( event.model_dump(), f"{event.type}", exchange=exchange, exchange_type="topic", persistent=True, headers={ "trace_id": str(uuid.uuid4()), "event_type": event.type, "schema_version": "1.0" }, correlation_id=event.id ) # Использование async def create_order(user_id: int, items: list): order_id = str(uuid.uuid4()) # Сохранить заказ await db.orders.insert({"id": order_id, "user_id": user_id}) # Опубликовать событие event = DomainEvent( id=str(uuid.uuid4()), type="order.created", aggregate_id=order_id, timestamp=datetime.utcnow().isoformat(), payload={"user_id": user_id, "items": items} ) await publish_event(event, "orders.events")

Разбор архитектурных решений:

#1. Domain Event как Pydantic модель

class DomainEvent(BaseModel): id: str # Уникальный ID события (для дедупликации) type: str # Тип события (order.created, user.signup) aggregate_id: str # ID сущности, которую изменили timestamp: str # Когда произошло событие (не когда опубликовано!) payload: dict # Бизнес-данные события

Почему не просто dict:

  • Валидация — Pydantic проверит поля до публикации
  • Документация — схема события явная и типизирована
  • AsyncAPI — автоматически сгенерируется документация
  • Эволюция схемы — легче добавить поля с defaults

#2. Topic exchange для событий

exchange_type="topic" # Позволяет routing по паттернам

Почему topic exchange:

  • Подписчики могут фильтровать события по паттерну: order.*, *.created, order.shipped
  • Гибче чем direct (точный match) и fanout (broadcast)

Production routing примеры:

# Подписчик всех событий заказов @broker.subscriber("order.*", exchange="orders.events", exchange_type="topic") async def handle_any_order_event(event: dict): ... # Подписчик только на creation @broker.subscriber("order.created", exchange="orders.events", exchange_type="topic") async def handle_order_created(event: dict): ... # Подписчик всех creation событий @broker.subscriber("*.created", exchange="orders.events", exchange_type="topic") async def handle_any_created(event: dict): ...

#3. Критичная проблема: нет гарантии доставки

Текущий код имеет anti-pattern:

async def create_order(user_id: int, items: list): order_id = str(uuid.uuid4()) # 1. Сохранить заказ await db.orders.insert({"id": order_id, "user_id": user_id}) # 2. Опубликовать событие event = DomainEvent(...) await publish_event(event, "orders.events") # ⚠️ Может упасть!

Если step 2 упадет: заказ в БД есть, но событие не опубликовано. Подписчики не узнают о заказе.

Решение 1: Outbox Pattern (production-ready)

async def create_order_with_outbox(user_id: int, items: list): order_id = str(uuid.uuid4()) event = DomainEvent( id=str(uuid.uuid4()), type="order.created", aggregate_id=order_id, timestamp=datetime.utcnow().isoformat(), payload={"user_id": user_id, "items": items} ) # Атомарно сохраняем заказ И событие в БД await db.transactional_save_order_and_event(order, event) # Отдельный процесс (CDC poller) прочитает outbox и опубликует # Это гарантирует доставку даже если broker down

Решение 2: Transactional outbox с Debezium

Application → DB (order + event in outbox table)
                   ↓
            Debezium CDC
                   ↓
              Kafka topic

Debezium читает DB binlog и публикует изменения в Kafka — атомарность на уровне БД.


#Best Practices для Production

#✅ DO: Используйте Pydantic модели для сообщений

class OrderCreated(BaseModel): order_id: str user_id: int total: float items: list[str] await broker.publish( OrderCreated(order_id="123", user_id=1, total=99.9, items=["item1"]), "orders.created" )

Почему: Валидация на отправке, автодокументирование, type safety.

#✅ DO: Всегда указывайте correlation_id для tracing

import uuid trace_id = str(uuid.uuid4()) await broker.publish( data, "channel", correlation_id=trace_id, headers={"trace_id": trace_id} )

Почему: Без tracing невозможно debug distributed systems.

#✅ DO: Используйте retry с exponential backoff

async def publish_with_retry(data, channel): for attempt in range(3): try: await broker.publish(data, channel, raise_timeout=True) return except PublishError: await asyncio.sleep(2 ** attempt) # 1s, 2s, 4s raise PublishFailed("All retries exhausted")

Почему: Сеть ненадежна, transient failures случаются часто.

#✅ DO: Monitor publish latency

import time start = time.monotonic() await broker.publish(data, channel) latency = time.monotonic() - start metrics.histogram("publish_latency_ms", latency * 1000)

Почему: Рост latency — ранний сигнал о проблемах брокера.

#✅ DO: Используйте batch для high throughput

async with RabbitBatch(broker, batch_size=100) as batch: for item in items: batch.publish(item, "channel")

Почему: 10x больше throughput при тех же ресурсах.


#❌ ANTI-PATTERNS: Как НЕ надо делать

#❌ DON'T: Публикация без обработки ошибок

# ПЛОХО — ошибка публикации теряется await broker.publish(data, "channel")

Проблема: Если broker down, сообщение теряется тихо.

Хорошо:

try: await broker.publish(data, "channel", raise_timeout=True) except PublishError as e: logger.error(f"Publish failed: {e}") # Retry или dead letter queue

#❌ DON'T: Публикация перед сохранением в БД

# ПЛОХО — сначала publish, потом DB await broker.publish({"order_id": 123}, "orders.created") await db.orders.insert(order) # Может упасть!

Проблема: Событие опубликовано, но заказа нет — инконсистентность.

Хорошо: Сначала БД, потом publish (или outbox pattern).

#❌ DON'T: Большие payload в сообщениях

# ПЛОХО — 10MB payload await broker.publish({ "order_id": 123, "huge_payload": giant_json_blob # 10MB }, "orders.created")

Проблема:

  • Broker memory pressure (каждое сообщение копируется)
  • Slow network transfer
  • Consumer OOM risk

Хорошо: Публикуйте ID, данные загруйте отдельно:

await broker.publish({ "order_id": 123, "data_url": f"https://storage/orders/123.json" # Ссылка на данные }, "orders.created")

#❌ DON'T: Один exchange для всех типов событий

# ПЛОХО — все события в один exchange await broker.publish(order_event, "events", exchange="all-events") await broker.publish(user_event, "events", exchange="all-events") await broker.publish(payment_event, "events", exchange="all-events")

Проблема:

  • Подписчики получают ненужные события (filter overhead)
  • Невозможно масштабировать отдельно
  • Blast radius при проблемах

Хорошо: Разделите по domain:

await broker.publish(order_event, "orders", exchange="orders.events") await broker.publish(user_event, "users", exchange="users.events")

#❌ DON'T: Игнорирование persistent flag для критичных данных

# ПЛОХО — заказы без persistent await broker.publish({"order_id": 123}, "orders.created") # persistent=False by default

Проблема: При restart RabbitMQ заказы потеряны.

Хорошо:

await broker.publish( {"order_id": 123}, "orders.created", persistent=True # Явно указываем )

#Decision Tree: Что выбрать?

Нужно опубликовать событие
    |
    ├── Критичные данные (заказы, платежи)?
    │   ├── Да → persistent=True + retry + outbox pattern
    │   └── Нет → persistent=False для скорости
    |
    ├── Высокий throughput (>10K msg/s)?
    │   ├── Да → Kafka или RabbitMQ batch
    │   └── Нет → RabbitMQ проще
    |
    ├── Нужен message replay?
    │   ├── Да → Kafka (хранит историю)
    │   └── Нет → RabbitMQ
    |
    ├── Ultra-low latency (<1ms)?
    │   ├── Да → NATS или Redis
    │   └── Нет → RabbitMQ/Kafka
    |
    └── Simple pub/sub без routing?
        ├── Да → Redis pub/sub
        └── Нет → RabbitMQ exchanges

Далее: Подписка и обработка сообщений (Subscribers)