Событийная архитектура: брокеры, топики, гарантии доставки
Синхронный вызов — это вопрос: «сделай это и скажи мне, что получилось». Событие — это утверждение: «вот что произошло; кому надо — разберётся». Разница выглядит косметической, но она меняет направление зависимости, модель отказа, модель согласованности и то, как команды в компании договариваются между собой. Событийная архитектура — не «микросервисы, но через очередь». Это отдельный способ мышления о системе, у которого есть своя цена, и цена эта платится в отладке, наблюдаемости и объяснении бизнесу, почему баланс обновляется «через секунду», а не «сразу».
Эта статья — про механику и про инженерную честность. Мы разберём, что такое событие в строгом смысле, чем лог отличается от очереди, что на самом деле означают три уровня гарантий доставки (и почему 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 показывает, что лог — это не «ещё одна очередь», а общий примитив, к которому сводятся репликация БД, стриминг и интеграция.
append-only, retention 7d")] T -->|offset 1042| G1[группа billing] T -->|offset 55| G2[группа analytics] T -->|offset 0 — перечитать всё| G3[новая группа search] end style T fill:none,stroke:#5aa469,stroke-width:2px style EX fill:none,stroke:#4a90d9,stroke-width:2px
Практический выбор:
| Нужно | Берите |
|---|---|
| Распределение задач воркерам, приоритеты, задержки, сложная маршрутизация | Очередь (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; брокер распределяет партиции между членами группы так, что каждая партиция читается ровно одним потребителем группы.
Отсюда следуют четыре факта, которые нужно выучить наизусть:
- Порядок гарантирован только внутри партиции. Глобального порядка в топике нет и быть не может — это цена горизонтального масштабирования.
- Ключ партиционирования — это ваша единица упорядоченности. Хотите, чтобы события одного заказа обрабатывались строго по порядку — ключ
order_id. Хотите порядок по клиенту — ключcustomer_id. Ключ выбирается один раз и меняется только через новый топик. - Число партиций — потолок параллелизма группы. Двенадцать партиций — максимум двенадцать полезных потребителей. Тринадцатый простаивает.
- Партиции легко добавить и невозможно убавить, и добавление ломает соответствие ключ→партиция для будущих сообщений (старые остаются где были). Планируйте с запасом ×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 — вопрос дизайна вашего обработчика.
Что делает 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 делается идемпотентностью на стороне эффекта, и никак иначе.
в одной транзакции БД Обработка --> ОшибкаВременная: таймаут, 503, дедлок Обработка --> ОшибкаПостоянная: невалидная схема,
нарушен инвариант ОшибкаВременная --> Обработка: retry с экспоненциальным
backoff и джиттером ОшибкаВременная --> DLQ: исчерпаны попытки ОшибкаПостоянная --> DLQ: без ретраев Применено --> Зафиксировано: commit offset Пропущено --> Зафиксировано: commit offset DLQ --> Зафиксировано: commit offset Зафиксировано --> [*]: следующее сообщение
Обратите внимание на структурное решение: временные и постоянные ошибки обрабатываются по-разному. Ретраить невалидную схему бессмысленно — это бесконечный цикл, который на 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. Полный поток: как это выглядит целиком
продукт должен показать промежуточный статус loop каждые 100 мс R->>DB: SELECT ... published_at IS NULL FOR UPDATE SKIP LOCKED R->>K: produce(key=order_id, OrderPlaced) K-->>R: ack (acks=all) R->>DB: UPDATE outbox SET published_at = now() end Note over R,K: падение между ack и UPDATE →
повторная публикация. Дубликат неизбежен. K-->>B: OrderPlaced (offset 1042) B->>B: inbox: event_id новый? списание + отметка — одна транзакция B->>K: commit offset 1042, produce(PaymentCaptured) K-->>N: OrderPlaced (та же запись, своя группа) N->>N: письмо (внешний эффект — идемпотентность по event_id!) K-->>O: PaymentCaptured O->>DB: UPDATE orders SET status='paid' Note over U,O: клиент увидит 'paid' при следующем опросе
или получит push — это и есть eventual consistency
Три вещи, на которые стоит смотреть в этой диаграмме. Во-первых, 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, ElixirPhoenix.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.
Источники
- Martin Fowler. What do you mean by “Event-Driven”? · Pat Helland. Life beyond Distributed Transactions
- Jay Kreps. The Log: What every software engineer should know about real-time data’s unifying abstraction
- Martin Kleppmann. Designing Data-Intensive Applications, гл. 11 «Stream Processing» — dataintensive.net
- Gregor Hohpe, Bobby Woolf. Enterprise Integration Patterns — enterpriseintegrationpatterns.com
- Tyler Treat. You Cannot Have Exactly-Once Delivery · KIP-98
- Chris Richardson. Pattern: Transactional outbox · Debezium Documentation
- Apache Kafka Documentation · RabbitMQ Quorum Queues · NATS JetStream
- CloudEvents Specification · OpenTelemetry Messaging Semantic Conventions
Что дальше
Мы научились надёжно передавать факты между сервисами. Следующий шаг — сделать эти факты основой модели данных внутри сервиса: разделить путь записи и путь чтения и хранить не текущее состояние, а поток событий, из которого оно выводится. Читайте «CQRS и Event Sourcing».