Архитектурные паттерны Событийная архитектура: брокеры, топики, гарантии доставки
0%

Событийная архитектура: брокеры, топики, гарантии доставки

Событийная архитектура: брокеры, топики, гарантии доставки

Синхронный вызов — это вопрос: «сделай это и скажи мне, что получилось». Событие — это утверждение: «вот что произошло; кому надо — разберётся». Разница выглядит косметической, но она меняет направление зависимости, модель отказа, модель согласованности и то, как команды в компании договариваются между собой. Событийная архитектура — не «микросервисы, но через очередь». Это отдельный способ мышления о системе, у которого есть своя цена, и цена эта платится в отладке, наблюдаемости и объяснении бизнесу, почему баланс обновляется «через секунду», а не «сразу».

Эта статья — про механику и про инженерную честность. Мы разберём, что такое событие в строгом смысле, чем лог отличается от очереди, что на самом деле означают три уровня гарантий доставки (и почему exactly-once в маркетинговом смысле не существует), как устроены партиции и порядок, как написать потребителя, который переживает дубликаты, и какие ошибки встречаются в каждой второй событийной системе.

Предполагается, что вы знакомы с материалом статей «Микросервисы: границы, коммуникация, данные, эксплуатация» и «Модульный монолит» — понятия «ограниченный контекст», «владелец данных» и «распределённый монолит» дальше используются без пояснений.


1. Что такое «событие» и почему это не просто сообщение

Событие — это неизменяемый факт о прошлом, принадлежащий своему источнику. Три слова здесь несут вес.

Факт о прошлом. OrderPlaced, а не PlaceOrder. Событие описывает то, что уже случилось и не может быть отменено (отмена — это новое событие OrderCancelled, а не удаление старого). Имена событий — всегда глагол в прошедшем времени.

Неизменяемый. Опубликованное событие нельзя отредактировать. Если в нём ошибка — публикуется компенсирующее или корректирующее событие. Это ровно та дисциплина, которая делает возможным воспроизведение истории.

Принадлежащий источнику. Событие описывает мир с точки зрения издателя и не знает, кто его прочитает. Как только в событии появляется поле send_email_to, оно перестало быть событием и стало командой в маскировке.

Различие «команда против события» — не педантизм, а вопрос о том, где живёт связность:

Команда Событие
Смысл «сделай Х» «случилось Y»
Имя повелительное наклонение прошедшее время
Получателей ровно один ноль или много
Кто знает о ком отправитель знает получателя издатель не знает подписчиков
Можно отклонить да, валидацией нет, факт уже произошёл
Транспорт обычно очередь point-to-point обычно топик pub/sub

Мартин Фаулер в статье «What do you mean by “Event-Driven”?» разбирает четыре разных паттерна, которые все называют одним словом, и это до сих пор лучшая карта местности:

Первые два — про интеграцию между сервисами и составляют 90% практики. Последние два — про внутреннее устройство сервиса и подробно разбираются в статье «CQRS и Event Sourcing».

Тонкое или толстое событие

Первое реальное проектное решение, и универсального ответа у него нет.

// Тонкое (event notification):
{"type": "order.placed", "order_id": "A-771", "occurred_at": "2026-07-16T10:00:00Z"}
// Толстое (event-carried state transfer):
{"type": "order.placed", "order_id": "A-771", "occurred_at": "...",
 "customer": {"id": "C-9", "tier": "gold"},
 "items": [{"sku": "X1", "qty": 2, "price_cents": 4990}],
 "total_cents": 9980, "currency": "RUB"}

Тонкое событие держит контракт крошечным, но каждый потребитель побежит обратно в сервис заказов за деталями — вы построили pub/sub поверх синхронного вызова и получили и связность, и асинхронность одновременно. Толстое событие устраняет обратные вызовы, но фиксирует куски вашей внутренней модели в публичном контракте и заставляет каждого потребителя думать об устаревании данных.

Практическое правило: публикуйте то, что нужно большинству потребителей, и ничего из того, что является вашей внутренней реализацией. Если три из пяти подписчиков сразу после события идут за одним и тем же полем — оно должно быть в событии. Пэт Хелланд формулирует это как разницу между «данными внутри» и «данными снаружи» в классической работе Data on the Outside versus Data on the Inside: данные снаружи неизменяемы, версионируемы и всегда немного из прошлого.

2. Очередь и лог: два разных примитива

Слово «брокер» скрывает две принципиально разные конструкции.

Очередь (broker-centric, RabbitMQ, SQS, ActiveMQ). Брокер хранит сообщения, отдаёт их потребителям и удаляет после подтверждения. Брокер знает состояние каждого сообщения: доставлено, в обработке, подтверждено, отклонено. Отсюда возможности: per-message ack/nack, приоритеты, TTL на сообщение, отложенная доставка, гибкая маршрутизация (в RabbitMQ — exchange с topic/fanout/direct/headers). Отсюда же ограничения: состояние на сообщение — это дорого, повторно прочитать историю невозможно (её нет), и очередь плохо переживает миллионы неподтверждённых сообщений.

Лог (log-centric, Kafka, Pulsar, NATS JetStream, Redpanda). Брокер — это append-only файл с индексом. Сообщение записывается в конец и живёт до истечения retention независимо от того, прочитал его кто-то или нет. Потребитель хранит собственную позицию (offset). Отсюда: несколько независимых групп потребителей читают одни и те же данные, историю можно перечитать сначала, пропускная способность упирается в последовательную запись на диск (сотни МБ/с на брокер). Отсюда же ограничения: нет per-message ack, порядок и параллелизм жёстко связаны через партиции, «отложить одно сообщение на час» — нетривиальная задача.

Джей Крепс в фундаментальном эссе The Log: What every software engineer should know about real-time data’s unifying abstraction показывает, что лог — это не «ещё одна очередь», а общий примитив, к которому сводятся репликация БД, стриминг и интеграция.

Практический выбор:

Нужно Берите
Распределение задач воркерам, приоритеты, задержки, сложная маршрутизация Очередь (RabbitMQ quorum queues, SQS)
Поток фактов, много независимых потребителей, перечитывание истории, стриминговая аналитика Лог (Kafka, Pulsar, JetStream)
Строгий FIFO на ключ при малом объёме, минимум эксплуатации SQS FIFO, Azure Service Bus sessions
Низкая латентность, лёгкий рантайм, request/reply вдобавок NATS / JetStream

Худший вариант — выбрать инструмент по популярности и потом воевать с его моделью. Kafka как «очередь задач с приоритетами» и RabbitMQ как «хранилище истории событий» — два самых частых способа потратить квартал.

3. Топики, партиции, порядок

Всё, что дальше, показано на примере Kafka, потому что её модель стала фактическим словарём отрасли; в Pulsar и JetStream названия другие, идеи те же.

Топик делится на партиции. Партиция — отдельный append-only лог со своей нумерацией (offset). Продюсер выбирает партицию: если задан ключ — hash(key) % partitions, если нет — по кругу (sticky round-robin). Потребители объединяются в consumer group; брокер распределяет партиции между членами группы так, что каждая партиция читается ровно одним потребителем группы.

Анатомия топика: партиции, ключи, offset и consumer group

Отсюда следуют четыре факта, которые нужно выучить наизусть:

  1. Порядок гарантирован только внутри партиции. Глобального порядка в топике нет и быть не может — это цена горизонтального масштабирования.
  2. Ключ партиционирования — это ваша единица упорядоченности. Хотите, чтобы события одного заказа обрабатывались строго по порядку — ключ order_id. Хотите порядок по клиенту — ключ customer_id. Ключ выбирается один раз и меняется только через новый топик.
  3. Число партиций — потолок параллелизма группы. Двенадцать партиций — максимум двенадцать полезных потребителей. Тринадцатый простаивает.
  4. Партиции легко добавить и невозможно убавить, и добавление ломает соответствие ключ→партиция для будущих сообщений (старые остаются где были). Планируйте с запасом ×2–×3 от текущей потребности, но не 1000 партиций «на всякий случай»: каждая партиция — это файловые дескрипторы, память на брокере и время восстановления.

Горячая партиция — самая частая эксплуатационная беда. Если 40% трафика приходится на одного клиента-гиганта, а ключ — customer_id, одна партиция получает 40% нагрузки и лаг растёт только у неё. Лечится либо составным ключом (customer_id:bucket, где bucket — хеш от order_id по модулю K, если внутри клиента строгий порядок не нужен), либо отдельным топиком для «китов».

Ребалансировка — процесс переназначения партиций при входе/выходе потребителя. Классический eager-протокол на время ребаланса останавливает всю группу («stop-the-world»); при частых деплоях это заметно. Современный ответ — кооперативный incremental rebalance (CooperativeStickyAssignor) и статическое членство (group.instance.id), которое переживает перезапуск пода без ребаланса. Детали — в документации Kafka по consumer configs.

4. Гарантии доставки: что можно обещать честно

Здесь начинается место, где маркетинг расходится с теорией.

Сеть может потерять запрос, может потерять ответ, и получатель не в силах отличить одно от другого. Это задача двух генералов, и она доказуемо неразрешима. Следствие: отправитель, не получивший подтверждения, обязан выбрать между «повторить» (рискуя дубликатом) и «не повторять» (рискуя потерей). Третьего варианта в распределённой системе нет.

  • At-most-once — подтверждаем до обработки. Потерь возможны, дубликатов нет. Годится для метрик, телеметрии, «последнее известное значение».
  • At-least-once — подтверждаем после обработки. Дубликаты возможны, потерь нет. Промышленный выбор по умолчанию.
  • Exactly-once — недостижим как свойство доставки, но достижим как свойство эффекта: сообщение может прийти дважды, но результат обработки применяется один раз. Этого добиваются идемпотентностью или транзакцией, охватывающей и эффект, и фиксацию позиции.

Тайлер Трит разбирает это подробно в You Cannot Have Exactly-Once Delivery, и вывод стоит повесить на стену: exactly-once delivery невозможна; exactly-once processing — вопрос дизайна вашего обработчика.

Окно потери и окно дубликата: когда фиксировать offset

Что делает Kafka EOS и где его границы

Kafka с версии 0.11 (KIP-98) умеет две вещи. Первая — идемпотентный продюсер (enable.idempotence=true): каждому продюсеру выдаётся PID, каждому сообщению — порядковый номер, и брокер отбрасывает повторы при ретрае. Это устраняет дубликаты, порождённые сетевым ретраем на участке producer→broker. Вторая — транзакции: атомарная запись в несколько партиций плюс атомарная фиксация offset потребления, то есть паттерн «прочитал–обработал–записал» (read-process-write) целиком. Проектный документ — KIP-98, разбор — в статье Confluent.

Критически важная оговорка: транзакции Kafka атомарны только внутри Kafka. Как только ваш обработчик пишет в PostgreSQL, дёргает платёжный API или отправляет письмо — гарантия заканчивается на границе. Для таких случаев exactly-once делается идемпотентностью на стороне эффекта, и никак иначе.

Обратите внимание на структурное решение: временные и постоянные ошибки обрабатываются по-разному. Ретраить невалидную схему бессмысленно — это бесконечный цикл, который на at-least-once потребителе останавливает всю партицию (классический poison pill, блокирующий лаг для всех последующих сообщений).

5. Идемпотентность на практике

Идемпотентный обработчик — тот, для которого $f(f(x)) = f(x)$: повторная обработка того же события не меняет результат. Есть три уровня.

Уровень 0. Естественная идемпотентность. SET status = 'shipped' идемпотентна сама по себе, balance = balance + 100 — нет. Если можно выразить эффект как присваивание, а не как инкремент, — сделайте так.

Уровень 1. Условное обновление. Версия/состояние в WHERE: UPDATE orders SET status='paid', version=version+1 WHERE id=? AND version=?. Повтор просто не находит строку.

Уровень 2. Таблица inbox (дедупликация по идентификатору). Универсальный способ: у каждого события есть уникальный event_id, потребитель хранит обработанные идентификаторы и в одной транзакции применяет эффект и записывает факт обработки.

-- Таблица дедупликации на стороне потребителя.
-- PRIMARY KEY даёт атомарную проверку «видели ли мы это событие».
CREATE TABLE inbox_processed (
    consumer     text        NOT NULL,   -- имя обработчика: одно событие могут
    event_id     uuid        NOT NULL,   -- обрабатывать несколько потребителей
    processed_at timestamptz NOT NULL DEFAULT now(),
    PRIMARY KEY (consumer, event_id)
);
-- Обязательно: чистка старых записей, иначе таблица растёт бесконечно.
CREATE INDEX ON inbox_processed (processed_at);
-- Retention inbox должен быть заметно больше retention топика,
-- иначе перечитывание истории породит повторное применение эффектов.
"""Идемпотентный потребитель Kafka: ручной коммит, inbox-дедупликация, DLQ."""
import json, logging, psycopg
from confluent_kafka import Consumer, Producer, KafkaException

CONSUMER_NAME = "billing.charge"

consumer = Consumer({
    "bootstrap.servers": "kafka:9092",
    "group.id": CONSUMER_NAME,
    "enable.auto.commit": False,         # коммитим руками, только после эффекта
    "auto.offset.reset": "earliest",     # новая группа читает историю с начала
    "isolation.level": "read_committed",  # не видим незакоммиченные транзакции
    "max.poll.interval.ms": 300_000,     # запас времени на медленную обработку
    "partition.assignment.strategy": "cooperative-sticky",
})
dlq = Producer({"bootstrap.servers": "kafka:9092", "enable.idempotence": True})


class PermanentError(Exception):
    """Ошибка, которую бессмысленно ретраить: битая схема, нарушенный инвариант."""


def handle(msg, conn) -> None:
    """Эффект и отметка о дедупликации — атомарно, в одной транзакции БД."""
    event = json.loads(msg.value())
    with conn.transaction():
        # Пытаемся застолбить event_id. Конфликт означает: уже обработано.
        cur = conn.execute(
            "INSERT INTO inbox_processed (consumer, event_id) VALUES (%s, %s) "
            "ON CONFLICT DO NOTHING RETURNING event_id",
            (CONSUMER_NAME, event["id"]),
        )
        if cur.fetchone() is None:
            logging.info("дубликат %s пропущен", event["id"])
            return
        # Атомарность «эффект + отметка» и есть exactly-once processing.
        conn.execute(
            "UPDATE accounts SET balance_cents = balance_cents - %s "
            "WHERE id = %s AND balance_cents >= %s",
            (event["amount_cents"], event["account_id"], event["amount_cents"]),
        )


def run() -> None:
    consumer.subscribe(["payments.v1"])
    with psycopg.connect("postgresql://app@db/app") as conn:
        while True:
            msg = consumer.poll(1.0)
            if msg is None:
                continue
            if msg.error():
                raise KafkaException(msg.error())
            try:
                handle(msg, conn)
            except PermanentError as exc:
                # Яд не должен блокировать партицию: уводим в DLQ и идём дальше.
                dlq.produce("payments.v1.dlq", msg.value(),
                            headers=[("error", str(exc).encode()),
                                     ("origin_offset", str(msg.offset()).encode())])
                dlq.flush()
            except Exception:
                # Временная ошибка: НЕ коммитим offset — сообщение придёт снова.
                logging.exception("временный сбой, offset не сдвигаем")
                consumer.seek(msg)
                continue
            consumer.commit(msg, asynchronous=False)  # только теперь фиксируем

Стоимость дедупликации: одна вставка с проверкой уникальности на событие, O(log n) по B-tree индексу, порядка десятков микросекунд на попадании в кэш. При 10 000 событий/с это ~10k доп. вставок/с — заметная, но обычно приемлемая нагрузка. Если она неприемлема, используют вероятностный фильтр (Bloom) как быстрый отсев перед точной проверкой либо дедупликацию в самом стриминговом движке (Flink с RocksDB-состоянием и чекпоинтами).

6. Dual write и transactional outbox

Самая коварная ошибка событийных систем занимает две строки кода:

# СЛОМАНО: два разных хранилища, одна транзакция невозможна
db.execute("UPDATE orders SET status='paid' WHERE id=%s", (order_id,))
db.commit()
producer.produce("orders.v1", event)   # ← упало? БД изменена, событие потеряно

Поменяв строки местами, вы поменяете тип бага: событие опубликовано, БД не изменена. Атомарности между PostgreSQL и Kafka нет; двухфазный коммит (XA) технически существует, но блокирует ресурсы на время неопределённости и в высоконагруженных системах практически не применяется — см. разбор Пэта Хелланда Life beyond Distributed Transactions.

Рабочее решение — transactional outbox (microservices.io): событие записывается в таблицу той же БД, в той же транзакции, что и изменение состояния. Отдельный процесс (релей) вычитывает таблицу и публикует в брокер, повторяя до успеха. Получается at-least-once с гарантией «событие есть тогда и только тогда, когда изменение зафиксировано».

-- Публикация события и изменение состояния — атомарно, в одной транзакции.
BEGIN;
  UPDATE orders SET status = 'paid', version = version + 1
   WHERE id = 'A-771' AND version = 7;
  INSERT INTO outbox (event_id, aggregate_id, topic, type, payload)
  VALUES (gen_random_uuid(), 'A-771', 'orders.v1', 'order.paid',
          '{"order_id":"A-771","amount_cents":9980}'::jsonb);
COMMIT;

-- Релей: забирает пачку, не мешая конкурентам (SKIP LOCKED), порядок по seq.
BEGIN;
SELECT seq, event_id, aggregate_id, topic, type, payload FROM outbox
 WHERE published_at IS NULL ORDER BY seq LIMIT 500 FOR UPDATE SKIP LOCKED;
-- ... publish в брокер с key = aggregate_id ...
UPDATE outbox SET published_at = now() WHERE seq = ANY(%s);
COMMIT;

Два способа реализовать релей:

  • Polling publisher — простой SELECT ... FOR UPDATE SKIP LOCKED в цикле. Плюс: нет новых компонентов, легко отлаживать. Минус: латентность = период опроса (обычно 50–500 мс), нагрузка на БД холостыми запросами.
  • Log tailing / CDC — читаем журнал репликации БД (Debezium поверх WAL PostgreSQL или binlog MySQL). Плюс: околонулевая задержка, нулевая нагрузка на таблицы, гарантированный порядок. Минус: ещё один распределённый компонент со своей эксплуатацией, чувствительность к схеме и к слотам репликации (незакрытый слот способен переполнить диск на мастере — самая частая авария Debezium).

Осторожно с очевидной ловушкой: не публикуйте наружу сырые CDC-события из таблиц. Это превращает вашу физическую схему в публичный контракт, и любая миграция колонки становится ломающим изменением для всех потребителей. CDC используют как транспорт для outbox-таблицы, спроектированной как контракт, а не как способ «стримить базу».

7. Схемы и эволюция контрактов

Событие переживёт код, который его породил. Через год ваш топик прочитает сервис, которого сейчас нет, написанный командой, которой сейчас нет. Поэтому схема — не формальность.

Реестр схем (Confluent Schema Registry, Apicurio) хранит версии схем и проверяет совместимость при регистрации. Ключевые режимы (документация Confluent):

Режим Что можно Кого обновляем первым
BACKWARD (по умолчанию) удалить поле, добавить поле со значением по умолчанию потребителей
FORWARD добавить поле, удалить поле с default издателей
FULL только добавление/удаление полей с default любого
NONE всё молитесь

Правило, которое покрывает 95% случаев: новые поля — только опциональные с дефолтом; поля не удаляются и не меняют тип; переименование = новое поле + период двойной записи. Ломающее изменение — это новый топик (orders.v2) и период параллельной публикации, а не «мы просто выкатим в понедельник».

Форматы: Avro (компактный, схема в реестре, отличная поддержка эволюции), Protobuf (строгий, хорошая кодогенерация, привычен из gRPC — см. «Стили API»), JSON Schema (читаемо и отлаживаемо, дороже по объёму). Для метаданных конверта имеет смысл стандарт CloudEvents — он фиксирует набор атрибутов (id, source, type, time, subject, datacontenttype) и маппинги на Kafka/HTTP/AMQP, избавляя от изобретения собственного конверта в каждой компании.

{
  "specversion": "1.0",
  "id": "6f1c0e02-9d4d-4a1a-9b0e-1c2b3a4d5e6f",
  "source": "/orders-service",
  "type": "com.acme.order.paid.v1",
  "subject": "A-771",
  "time": "2026-07-16T10:00:00Z",
  "datacontenttype": "application/json",
  "traceparent": "00-4bf92f3577b34da6a3ce929d0e0e4736-00f067aa0ba902b7-01",
  "data": {"order_id": "A-771", "amount_cents": 9980, "currency": "RUB"}
}

Поле traceparent — не украшение. В событийной системе стек вызовов не существует, и единственный способ ответить на вопрос «почему клиенту не пришло письмо» — сквозной контекст трассировки, пробрасываемый через заголовки сообщений (OpenTelemetry messaging conventions).

8. Полный поток: как это выглядит целиком

Три вещи, на которые стоит смотреть в этой диаграмме. Во-первых, 202 Accepted вместо 200 OK — событийная архитектура протекает в API и в UX, и делать вид, что это чисто внутреннее решение, нельзя. Во-вторых, дубликат при публикации структурно неизбежен (шаг между ack и UPDATE), поэтому дедупликация у потребителя — обязательный, а не опциональный элемент. В-третьих, две группы потребителей читают одну и ту же запись независимо: добавление сервиса уведомлений не потребовало ни строчки изменений в сервисе заказов. Ради этого свойства всё и затевалось.

9. Типичные ошибки

Событие-команда. UserRegisteredSendWelcomeEmail или поле next_action в payload. Симптом: издатель знает, что должно случиться дальше. Лечение: событие описывает только факт; решение о письме принимает подписчик.

Событийная связность вместо развязки. Сервис A публикует событие, ждёт события-ответа от B и без него не может продолжить. Это синхронный RPC, реализованный самым дорогим способом. Если нужен ответ здесь и сейчас — делайте честный синхронный вызов.

Предположение о глобальном порядке. «События же придут по порядку» — не придут, если у них разные ключи. Проектируйте обработчики так, чтобы они переживали переупорядочивание: версия/timestamp в событии и правило «применяю, только если версия больше текущей» (last-writer-wins по версии агрегата, а не по времени прибытия).

Брокер как база данных. «Retention бесконечный, состояние восстановим перечитыванием топика». Работает до первого инцидента, когда перечитывание 400 млн событий занимает шесть часов и роняет downstream. Лог — транспорт и журнал; долговременное состояние живёт в БД или в компактящемся топике с осознанным ключом.

DLQ без процесса. Dead letter queue есть, алерта на её глубину нет, и через полгода там 90 000 сообщений, которые никто не разберёт. DLQ без owner’а, дашборда и процедуры повторной подачи — это ритуал, а не механизм. Туда же: ретраи без backoff и джиттера — downstream упал, тысяча потребителей долбит его каждые 100 мс, и он не встаёт никогда (см. «Устойчивость: circuit breaker, retry, bulkhead, backpressure»).

Отсутствующий журнал причин. Никто не может ответить, почему у клиента баланс −300. В синхронной системе есть трейс; в событийной без сквозного trace-id и без сохранения событий вы отлаживаете по логам четырёх сервисов вручную.

Eventual consistency, не проговорённая с продуктом. Инженерное решение «баланс сойдётся через 800 мс» превращается в баг-репорт «деньги списались, а заказ не оплачен». Расхождение должно быть либо невидимым для пользователя (оптимистичный UI), либо явно показанным («обрабатывается»).

10. Эксплуатация: что настраивать и на что смотреть

Настройки долговечности (Kafka; в других брокерах есть аналоги):

# Продюсер: не терять данные при отказе брокера
acks: all                       # ждать подтверждения от всех ISR-реплик
enable.idempotence: true        # дедупликация ретраев на стороне брокера
max.in.flight.requests.per.connection: 5  # при идемпотентности порядок сохранён
delivery.timeout.ms: 120000     # ограничиваем временем, а не числом попыток
compression.type: zstd          # 3–5× экономии трафика и диска на JSON
linger.ms: 10                   # микробатчинг: +10 мс латентности, кратно больше throughput

# Топик: сколько отказов переживаем
replication.factor: 3
min.insync.replicas: 2          # с acks=all переживаем отказ одного брокера
unclean.leader.election.enable: false  # НИКОГДА true: молчаливая потеря данных
retention.ms: 604800000         # 7 суток — окно на восстановление потребителя

Связка acks=all + replication.factor=3 + min.insync.replicas=2 — канонический безопасный набор: запись подтверждается, только когда она есть минимум на двух брокерах, и отказ одного не останавливает приём. min.insync.replicas=3 при RF=3 выглядит «надёжнее», но останавливает запись при потере любого брокера — типичная ошибка перестраховки.

Метрики, на которые ставят алерты:

Метрика Почему важна Ориентир
Consumer lag (в сообщениях и в секундах) главный индикатор здоровья: растёт — не успеваем алерт на устойчивый рост, а не на абсолютное число
Lag по партициям по отдельности ловит горячую партицию и залипшего потребителя перекос > 5× между партициями
Глубина и приток DLQ тихая деградация любой ненулевой приток заслуживает разбора
Время обработки сообщения (p99) подбирается к max.poll.interval.ms → ребалансы p99 < 30% от лимита
Частота ребалансировок «дребезг» группы съедает пропускную способность > 1 в 10 минут — расследовать
Under-replicated partitions риск потери данных должно быть 0
Возраст самого старого неопубликованного outbox-события релей встал > 60 с — алерт

Consumer lag в секундах (а не в сообщениях) — куда более честная метрика: «300 000 сообщений отставания» ничего не значит, а «отстаём на 4 минуты» немедленно переводится в разговор с продуктом.

Порядки величин. Одна партиция Kafka на обычном железе — примерно 10–50 МБ/с и десятки тысяч сообщений в секунду; латентность end-to-end при linger.ms=0 и acks=all — единицы миллисекунд внутри ДЦ. Число партиций считают от целевой пропускной способности: partitions ≈ max(target_throughput / per_partition_throughput, target_throughput / per_consumer_throughput), округляя вверх с запасом. RabbitMQ classic-очередь — десятки тысяч сообщений/с при небольшой глубине; quorum-очередь (документация) медленнее, но переживает отказ узла без потерь. Эти цифры — не обещание, а порядок для прикидки перед нагрузочным тестом.

Компактящиеся топики. cleanup.policy=compact заставляет Kafka хранить последнее значение для каждого ключа бесконечно, а старые версии удалять. Это превращает топик в реплицируемую таблицу «текущее состояние по ключу» — основа для локальных кэшей и материализованных представлений. Отправка null в качестве значения (tombstone) означает удаление ключа — именно так реализуют «право на забвение» в логе, из которого, казалось бы, ничего нельзя удалить.

11. Когда событийная архитектура не нужна

  • Нужен немедленный ответ с результатом. Проверка промокода, авторизация, поиск. Асинхронность здесь — это 202 Accepted и опрос, то есть худший UX за большие деньги. Туда же: один потребитель, который всегда будет один — прямой вызов проще и отлаживается стеком.
  • Строгий инвариант между двумя сущностями. «Сумма позиций всегда равна итогу заказа» внутри одного агрегата — это транзакция БД, а не сага. См. «Saga, распределённые транзакции, outbox и идемпотентность».
  • Команда меньше уровня зрелости. Событийная система требует трассировки, схем, DLQ-процесса, дежурства. Без этого она превращается в «сообщения куда-то уходят» — состояние, из которого выходят месяцами.
  • Внутри модульного монолита. Внутрипроцессная шина событий (Spring ApplicationEventPublisher, MediatR, Elixir Phoenix.PubSub) даёт развязку кода без сетевых отказов и часто это ровно то, что нужно. Брокер добавляют, когда появляется вторая единица развёртывания.

Разумный маршрут внедрения: сначала внутрипроцессные события с явными контрактами, затем outbox и один топик для одного реального сценария интеграции, затем — по мере появления потребителей — расширение. Обратный порядок («сначала поставим Kafka, потом придумаем, что в неё писать») даёт кластер, три топика и всю сложность без единой выгоды.

Мини-итог

  • Событие — неизменяемый факт о прошлом, принадлежащий издателю; команда адресована конкретному получателю. Смешение этих понятий — источник большинства проблем проектирования.
  • Очередь и лог — разные примитивы. Очередь хранит состояние на сообщение и умеет сложную маршрутизацию; лог хранит позицию у читателя и позволяет многим независимым потребителям перечитывать историю.
  • Порядок гарантирован только внутри партиции, ключ партиционирования = единица упорядоченности, число партиций = потолок параллелизма.
  • Exactly-once delivery не существует. Существует at-least-once + идемпотентный обработчик, что даёт exactly-once effect. Атомарность «эффект + отметка о дедупликации» — сердце корректного потребителя. Dual write ломается всегда, лечится transactional outbox (poll или CDC).
  • Схема — публичный контракт, живущий дольше кода. Обратная совместимость, реестр схем, версия в имени топика, новый топик вместо ломающего изменения.
  • В проде смотрят на lag в секундах, перекос по партициям, приток в DLQ и частоту ребалансов; из настроек критичны acks=all, min.insync.replicas=2, unclean.leader.election.enable=false.

Источники

Что дальше

Мы научились надёжно передавать факты между сервисами. Следующий шаг — сделать эти факты основой модели данных внутри сервиса: разделить путь записи и путь чтения и хранить не текущее состояние, а поток событий, из которого оно выводится. Читайте «CQRS и Event Sourcing».

Нашли неточность? Выделите фрагмент текста — рядом появится жучок.

Нужен разбор именно вашей ситуации?

Статья описывает общий случай. Если у вас частный — можно разобрать его отдельно, платно. А если не хватает целого материала, предложите тему: её оплачивают вскладчину, и она выходит открытой для всех.

Доска запросов