Очереди и потоки: брокеры, порядок сообщений, backpressure, повторы
Брокер сообщений выглядит как труба: с одной стороны положили, с другой достали. Это самая дорогая метафора в инженерии распределённых систем. Труба не имеет состояния, не выбирает лидера, не теряет данные при переключении и не может «переставить местами» воду. Брокер делает всё перечисленное.
Правильная метафора другая: брокер — это ещё одна распределённая база данных, которую вы поставили в самый горячий путь своей системы. У неё есть репликация (см. репликацию), выборы лидера (см. консенсус), модель согласованности, окно долговечности и собственный набор режимов отказа. Разница с обычной БД в том, что её отказы проявляются не как ошибка запроса, а как задержка в шесть часов, дубликат платежа или переставленные местами обновления профиля — то есть как баг бизнес-логики.
Эта статья — про четыре вещи, которые определяют, будет ли ваша событийная архитектура работать: какую модель брокера вы выбрали, что вы понимаете под порядком, как система ведёт себя при перегрузке и что происходит с сообщением, которое не удалось обработать. Стили вроде event-driven и saga разбираются в архитектурных паттернах; здесь — фундамент под ними.
Что покупает брокер и чем вы за это платите
Синхронный вызов A → B требует, чтобы B был жив, быстр и вмещал текущий трафик A. Брокер снимает все три требования сразу — и берёт за это плату.
| Что покупаете | Как это работает | Чем платите |
|---|---|---|
| Временная развязка | получатель может быть выключен час — сообщения ждут | нет обратной связи: отправитель не знает, что произошло с его событием |
| Сглаживание пиков | буфер поглощает всплеск на 10× | пик превращается в лаг, а не исчезает; SLA «онлайн» становится ложью |
| Fan-out | один producer, N независимых потребителей | N независимых наборов лагов, ретраев и DLQ |
| Развязка деплоя | схема сообщений вместо синхронного контракта | эволюция схемы становится отдельной инженерной дисциплиной |
| Сглаживание отказов | падение B не роняет A |
падение B теперь невидимо для A: нужен отдельный мониторинг |
Сценарий отказа №0, самый частый. Команда добавила очередь между оформлением заказа и списанием со склада, «чтобы убрать нагрузку с БД». Нагрузка не исчезла — она стала отложенной. В чёрную пятницу входной поток вырос в 8 раз, обработчик — нет. Через два часа лаг составил 40 минут, через шесть — 4 часа. Пользователь видит «Заказ принят», письмо приходит ночью, товар за это время уже продан кому-то ещё. В логах брокера — идеальная тишина: он справился, он всё сохранил. Проблема видна только на графике consumer_lag_seconds, которого никто не сделал.
Правило, из которого стоит исходить: очередь не увеличивает пропускную способность системы, она перераспределяет её во времени. Если устойчиво λ > μ (приходит больше, чем обрабатывается), очередь только откладывает отказ и делает его хуже.
Две модели: очередь и лог
Все брокеры делятся на два семейства по одному признаку — кто хранит позицию чтения.
В очереди (RabbitMQ, SQS, ActiveMQ) брокер отслеживает каждое сообщение индивидуально: выдал → пометил как in-flight → получил ack → удалил. Состояние на брокере пропорционально числу неподтверждённых сообщений. Чтение деструктивно: прочитанное исчезает.
В логе (Kafka, Pulsar, Kinesis, NATS JetStream, Redpanda) брокер хранит упорядоченную неизменяемую последовательность и вообще не знает, кто что прочитал. Позиция — это число (оффсет), которое коммитит сам консьюмер. Состояние на брокере пропорционально числу партиций, а не сообщений. Именно поэтому Kafka спокойно держит терабайты при неизменной нагрузке на память.
| Свойство | Очередь | Лог |
|---|---|---|
| Кто хранит позицию | брокер, per-message | консьюмер, per-partition |
| Чтение | деструктивное | курсорное, неразрушающее |
| Replay / backfill | невозможен без внешней копии | штатная операция: сдвинуть оффсет |
| Новый потребитель исторических данных | нужно перезалить | подписался и прочитал с начала |
| Порядок | FIFO только при одном потребителе | тотальный внутри партиции |
| Единица параллелизма | сообщение | партиция |
| Приоритеты, per-message TTL, сложная маршрутизация | из коробки | нет, делается топиками и приложением |
| Стоимость подтверждения | запись состояния на брокере | инкремент числа |
| Типичный масштаб | десятки тысяч очередей | десятки тысяч партиций |
Различие оформилось как «умный брокер и простой потребитель» против «простого брокера и умного потребителя» — формулировка из design-документа Kafka и из Enterprise Integration Patterns Хопа и Вульфа. Манифест лог-модели — статья Джея Крепса «The Log: What every software engineer should know about real-time data’s unifying abstraction»; там же объясняется, почему лог — это та же абстракция, что и WAL базы данных и репликация конечных автоматов.
Сценарий отказа для очереди. Классические зеркалируемые очереди RabbitMQ теряли данные при разделении сети: после pause_minority или ручного слияния «победившая» половина отбрасывала подтверждённые публикации. Это подробно измерено в отчёте Jepsen по RabbitMQ. Современный ответ — quorum queues на Raft, где подтверждение означает запись большинством; цена — заметно меньшая пропускная способность и больше диска.
Сценарий отказа для лога. Консьюмер лежал 8 часов, retention — 6 часов. Он возвращается и запрашивает оффсет, которого уже нет:
[2026-03-14 04:22:19,551] INFO [Consumer clientId=orders-2, groupId=orders]
Fetch position FetchPosition{offset=88133521} is out of range for partition orders-7,
resetting offset (org.apache.kafka.clients.consumer.internals.SubscriptionState)
[2026-03-14 04:22:19,553] INFO [Consumer clientId=orders-2, groupId=orders]
Resetting offset for partition orders-7 to position FetchPosition{offset=91204418}
Две строки INFO. Между ними — 3 070 897 молча пропущенных сообщений, потому что auto.offset.reset=latest. Поставьте earliest, если данные важнее скорости восстановления, и обязательно алерт на сам факт resetting offset: это событие уровня инцидента, а не INFO.
Долговечность: что означает «брокер принял»
Между «клиент получил ack» и «данные переживут отказ узла» лежат три независимых решения, и каждое из них по умолчанию настроено в сторону скорости.
- Сколько реплик подтвердили. Kafka:
acks=0 | 1 | all, причёмallозначает «все из текущего ISR», который может сжаться до одного узла — поэтомуmin.insync.replicasобязателен (подробно в репликации). RabbitMQ: publisher confirms плюс quorum queue. SQS: репликация по нескольким AZ выполняется всегда. - Дошло ли до диска. Kafka по умолчанию полагается на page cache и репликацию, а не на fsync каждой записи (
flush.messages,flush.ms). Это осознанный выбор: три копии в памяти трёх машин надёжнее одной копии на одном диске — но только если машины действительно независимы (стойка, AZ, питание). Коррелированный отказ ломает эту логику; см. модели отказов. - Подтвердил ли брокер отправителю. Продюсер без
publisher confirmsв AMQP или сacks=0в Kafka считает отправкой запись в TCP-сокет. Это буквально «отправил в никуда, надеюсь на лучшее».
Сценарий отказа. Сервис публикует в RabbitMQ через basic_publish без confirms и без mandatory. Опечатка в routing key — сообщение не попало ни в одну очередь и тихо отброшено брокером. Ни ошибки, ни лога, ни метрики. Обнаружено через три недели по расхождению отчётов. Лечится включением confirms и mandatory с обработчиком basic.return, а на стороне Kafka — acks=all и алертом на record-error-rate.
Порядок сообщений: чего у вас нет и что вместо этого есть
Первое, что нужно принять: глобального порядка событий в распределённой системе не существует. Это не ограничение брокеров, а следствие относительности одновременности — см. время и логические часы. Любой брокер даёт порядок только относительно чего-то конкретного.
| Система | Единица порядка | Что ломает |
|---|---|---|
| Kafka | партиция | смена числа партиций, ретраи без идемпотентного продюсера, параллельный консьюмер |
| RabbitMQ | очередь при одном consumer | prefetch > 1, несколько consumer’ов, requeue после nack |
| SQS standard | ничего, best effort | сама модель |
| SQS FIFO | MessageGroupId |
лимит 300 операций/с на группу без батчинга |
| Kinesis | шард | resharding |
| Pulsar | партиция; в режиме Key_Shared — ключ |
смена подписки, изменение числа консьюмеров |
| NATS JetStream | stream | несколько consumer’ов с AckPolicy explicit |
Ключ сообщения — это одновременно единица порядка и единица параллелизма, и в этом главный компромисс. Крупный ключ (tenant_id) даёт широкие гарантии и горячие партиции (см. партиционирование). Мелкий ключ (event_id) даёт идеальное распределение и никакого порядка. Правильный ключ — тот, внутри которого события действительно причинно связаны: обычно это идентификатор агрегата (order_id, account_id).
Шесть мест, где порядок ломается
Отдельно стоит подчеркнуть третью группу: порядок доставки не равен порядку применения. Брокер может доставить идеально упорядоченно, а консьюмер с ThreadPoolExecutor применит вразнобой — и это самый частый источник «невоспроизводимых» багов.
Сценарий отказа. Пользователь меняет email: сначала a@x.com, через 200 мс b@x.com. Оба события в одной партиции, порядок доставки правильный. Консьюмер раскладывает пачку по 16 потокам; поток с первым событием попал в GC-паузу на 300 мс. Итог — в базе a@x.com. В логах нет ни одной ошибки: обе записи прошли, обе вернули 200 OK. Единственный след — жалоба пользователя через неделю.
Лечится это тремя способами, и выбирать нужно осознанно:
- Последовательно внутри ключа, параллельно между ключами. Стандартное решение: диспетчеризация по
hash(key) % Nна N воркеров с собственными очередями. - Монотонность вместо порядка. У события есть версия, применяется только более новая:
UPDATE ... WHERE version < :v. Тогда порядок вообще перестаёт быть требованием — подробнее про фенсинг и версии в идемпотентности. - Коммутативные операции.
INCREMENTвместоSET, CRDT-подобные структуры (см. репликацию). Порядок не важен по построению — самый надёжный вариант, если предметная область позволяет.
import hashlib
from queue import Queue
from threading import Thread
class KeyedExecutor:
"""Параллелизм между ключами, строгая последовательность внутри ключа.
Инвариант: все сообщения с одним ключом всегда попадают в один и тот же воркер,
а внутри воркера обрабатываются строго по одному. Это восстанавливает порядок,
потерянный при простом пуле потоков.
Сложность: O(1) на диспетчеризацию, память O(N * queue_maxsize).
"""
def __init__(self, handler, workers: int = 16, queue_size: int = 256):
self.handler = handler
# Ограниченные очереди — это и есть backpressure: put() заблокируется,
# когда воркер не успевает, и poll() у консьюмера остановится сам.
self.queues = [Queue(maxsize=queue_size) for _ in range(workers)]
for q in self.queues:
Thread(target=self._loop, args=(q,), daemon=True).start()
def _slot(self, key: bytes) -> int:
# Не встроенный hash(): он рандомизируется между процессами (PYTHONHASHSEED),
# и после рестарта тот же ключ уехал бы к другому воркеру.
digest = hashlib.blake2b(key, digest_size=8).digest()
return int.from_bytes(digest, "big") % len(self.queues)
def submit(self, key: bytes, msg) -> None:
self.queues[self._slot(key)].put(msg) # блокирует при заполнении — так и надо
def _loop(self, q: Queue) -> None:
while True:
msg = q.get()
try:
self.handler(msg)
finally:
q.task_done()
def drain(self) -> None:
"""Дождаться завершения перед коммитом оффсета: иначе получим at-most-once."""
for q in self.queues:
q.join()
Ключевая деталь в конце: drain() перед коммитом оффсета. Без него коммит происходит раньше обработки — и падение процесса приводит к тихой потере, а не к повтору. Про правильный порядок «эффект, потом коммит» — в гарантиях доставки.
Head-of-line blocking
В логе порядок и блокировка — две стороны одной медали. Партиция упорядочена, значит сообщение №5 нельзя пропустить, чтобы взяться за №6. Одно «тяжёлое» или «ядовитое» сообщение останавливает всю партицию.
Сценарий отказа. В топике payments одно сообщение вызывает NullPointerException в парсере из-за поля, которого не было в схеме. Консьюмер не коммитит оффсет, перечитывает то же сообщение, снова падает. Процесс жив, health-check зелёный, CPU 100%, лаг растёт линейно:
02:14:31 ERROR [payments-3] handler failed offset=88133612 attempt=1 NullPointerException
02:14:31 ERROR [payments-3] handler failed offset=88133612 attempt=2 NullPointerException
02:14:31 ERROR [payments-3] handler failed offset=88133612 attempt=3 NullPointerException
... 41 000 одинаковых строк за минуту ...
Одинаковый offset в повторяющихся строках — сигнатура poison message. Отдельная метрика, которая ловит это мгновенно: max(offset) - min(committed_offset) не меняется, а error_rate растёт.
Четыре способа не блокировать партицию: вынести сбойное сообщение в retry-топик и продолжить; использовать Key_Shared в Pulsar (порядок сохраняется по ключу, но медленный ключ не блокирует остальные); развести тяжёлые и лёгкие события по разным топикам; ограничить число попыток и отправить в DLQ. Канонический разбор с продакшн-масштабом — «Building Reliable Reprocessing and Dead Letter Queues with Apache Kafka» от Uber.
Backpressure: почему буфер не решает, а откладывает
Backpressure — это не «дропать лишнее». Это механизм, которым потребитель сообщает производителю: сбавь темп. Разница принципиальна: дроп — локальное решение о потере данных, обратное давление — распространение информации о перегрузке вверх по цепи.
Математика здесь простая и безжалостная. Закон Литтла связывает среднее число заявок в системе, интенсивность потока и время пребывания:
$$L = \lambda W$$
Если приходит λ = 1000 сообщений/с, а обработка одного занимает W = 50 мс, то в системе в среднем L = 50 сообщений. Пока обслуживающая мощность μ больше λ, всё стабильно. Как только λ превышает μ хотя бы на 5%, очередь растёт линейно во времени и не стабилизируется никогда — буфер любого размера будет исчерпан, вопрос только в том, когда.
Второй факт: время ожидания растёт нелинейно с загрузкой. Для простейшей модели M/M/1 при коэффициенте использования ρ = λ/μ
$$W = \frac{1}{\mu - \lambda} = \frac{1}{\mu} \cdot \frac{1}{1 - \rho}$$
При ρ = 0,5 задержка вдвое выше идеальной, при ρ = 0,9 — в десять раз, при ρ = 0,99 — в сто. Отсюда практический вывод: запас мощности в 30–40% — это не расточительство, а способ держать хвостовые задержки конечными. Планирование ёмкости «под 95% утилизации» гарантирует, что p99 будет катастрофическим.
Третий факт, который обычно упускают: большой буфер ухудшает ситуацию. Он не увеличивает пропускную способность (μ не изменилась), зато увеличивает время пребывания и делает данные протухшими к моменту обработки. Это ровно тот же эффект, что bufferbloat в сетях — см. Getty, Nichols, «Bufferbloat: Dark Buffers in the Internet» (ACM Queue, 2011).
Механизмы обратного давления
| Механизм | Где встречается | Как работает | Ограничение |
|---|---|---|---|
| Pull-модель | Kafka, Pulsar, JetStream | консьюмер сам решает, когда и сколько взять | давит только на себя, продюсер продолжает писать |
| Кредиты | AMQP basic.qos, RabbitMQ credit flow, HTTP/2 windows, gRPC |
получатель выдаёт разрешение на N единиц | нужна поддержка с обеих сторон |
request(n) |
Reactive Streams, Project Reactor, Akka Streams | подписчик явно запрашивает объём | заражает весь стек асинхронностью |
| Блокировка | ограниченная очередь, buffer.memory + max.block.ms |
producer тормозит и упирается в TCP | давление доходит до пользователя — это надо решить заранее |
| Load shedding | rate limiter, 429/503, приоритеты | лишнее отбрасывается осознанно | нужен явный критерий, что можно терять |
| Спил на диск | Kafka по построению, tiered storage Pulsar | буфер большой и дешёвый | всё равно конечен, см. retention |
Спецификация Reactive Streams стоит того, чтобы её прочитать целиком (это буквально несколько страниц): она формализует правило, к которому все приходят опытным путём — «производитель не имеет права отправить больше, чем потребитель запросил».
отбрасываем по приоритету, отдаём 429"] B -->|"Нет: платежи, заказы"| D{"Есть ли куда копить?"} D -->|"Да: брокер с диском и retention"| E["Копим в брокере
консьюмер тормозит сам, лаг под алертом"] D -->|"Нет: буфер только в памяти"| F{"Можно ли замедлить источник?"} F -->|"Да: свой сервис, свой клиент"| G["Блокирующий backpressure
ограниченная очередь + кредиты"] F -->|"Нет: внешний трафик"| H["Приём с ограничением скорости
+ явный отказ вместо тихой деградации"] E --> I{"Лаг растёт устойчиво?"} I -->|"Да"| J["Масштабировать консьюмеров
до числа партиций, дальше — партиции"] I -->|"Нет, это пик"| K["Ничего не делать: буфер работает как задумано"] C --> L["Обязательно: метрика отброшенного"] G --> M["Обязательно: тайм-аут на блокировку,
иначе дедлок вместо деградации"] style F fill:#c2413f,fill-opacity:0.14,stroke:#c2413f style E fill:#3f9d6b,fill-opacity:0.18,stroke:#3f9d6b
Что видно в логах, когда backpressure работает и когда его нет
Работает. RabbitMQ достиг порога памяти и заблокировал публикаторов:
=INFO REPORT==== 14-Mar-2026::02:31:07 ===
vm_memory_high_watermark set. memory_high_watermark = 0.6
memory used = 6.4GB allowed = 6.0GB
=WARNING REPORT==== 14-Mar-2026::02:31:07 ===
memory resource limit alarm set on node rabbit@mq-2
**********************************************************
*** Publishers will be blocked until this alarm clears ***
Это не отказ, а работающая защита. Плохо здесь только одно: если клиент не обрабатывает connection.blocked, он просто зависнет без объяснения. Продюсеры обязаны подписываться на это уведомление и отдавать наверх осмысленную ошибку.
Не работает. Kafka-продюсер упёрся в buffer.memory, брокер медленный, батчи протухают:
org.apache.kafka.common.errors.TimeoutException: Expiring 143 record(s)
for orders-3:120000 ms has passed since batch creation
org.apache.kafka.common.errors.TimeoutException: Failed to allocate memory
within the configured max blocking time 60000 ms.
Второе исключение означает, что HTTP-хендлеры вашего сервиса стоят по 60 секунд каждый. Пул потоков веб-сервера исчерпан, балансировщик отдаёт 504 — при том, что «упала» только асинхронная ветка. Обратное давление дошло до пользователя, и это правильно с точки зрения физики, но должно быть осознанным решением: либо блокируем, либо отвечаем 202 и теряем событие с метрикой, третьего нет.
Практика настройки
- Prefetch (
basic.qos) в RabbitMQ. Слишком маленький — консьюмер простаивает на round-trip; слишком большой — сообщения скапливаются у одного потребителя, распределение становится неравномерным, а лиза удерживается зря. Ориентир:prefetch ≈ время round-trip / время обработки, округлённое вверх, но не больше нескольких сотен. Для медленных обработчиков (сотни мс) правильное значение часто равно 1–5. max.poll.recordsиmax.poll.interval.msв Kafka. Первое ограничивает размер пачки, второе — сколько вам разрешено её обрабатывать. Классическая ошибка: 500 записей по 100 мс = 50 секунд при лимите в 300 секунд — норма; но если обработка деградировала до 700 мс, получится 350 секунд, ребаланс и повторная обработка всей пачки другим узлом.pause()/resume(). Если внутренняя очередь заполнилась, правильная реакция —consumer.pause(partitions), а не «продолжать poll и складывать в память». Poll обязан вызываться (иначе heartbeat умрёт), но возвращать он должен ноль записей.- Тайм-аут на любую блокировку. Ограниченная очередь без тайм-аута на
put()превращает деградацию в дедлок.offer(msg, 5, SECONDS)с явной веткой «не влезло» — обязательно.
Повторы: где они живут в каждом брокере
At-least-once означает повторы, а повторы — это усилитель нагрузки, включающийся ровно в момент деградации. Общая механика бэкоффа, джиттера и бюджета ретраев разобрана в гарантиях доставки; здесь — брокерная часть, которая у всех устроена по-разному.
| Брокер | Счётчик попыток | Задержка перед повтором | DLQ |
|---|---|---|---|
| RabbitMQ (quorum) | x-delivery-count в заголовке |
TTL очереди + DLX или плагин delayed-message | delivery-limit → dead-letter exchange |
| SQS | ApproximateReceiveCount |
VisibilityTimeout, DelaySeconds до 15 мин |
redrive policy: maxReceiveCount → DLQ |
| Kafka | нет, считает приложение | нет; retry-топики с разными задержками | отдельный топик, руками |
| Pulsar | reconsumeLater, счётчик в retry letter topic |
задаётся при reconsumeLater |
deadLetterPolicy.maxRedeliverCount |
| NATS JetStream | num_delivered в метаданных |
AckWait + массив backoff |
MaxDeliver → advisory + termination |
У Kafka встроенных ретраев нет принципиально: сдвиг оффсета вперёд ради «повторим позже» противоречит модели лога. Отсюда паттерн retry-топиков — orders.retry.5s, orders.retry.1m, orders.retry.10m, orders.dlq. Обработчик, получив повторяемую ошибку, публикует сообщение в следующий retry-топик и коммитит оффсет исходного. Партиция не блокируется, задержка достигается тем, что консьюмер retry-топика перед обработкой ждёт до timestamp + delay (или просто pause, пока не наступит время). В Spring Kafka это @RetryableTopic.
Жизненный цикл сообщения с повторами
порядок нарушен InFlight --> Ready: лиза истекла
обработка ещё идёт — источник дублей InFlight --> Delayed: повторяемая ошибка,
попытка N из K Delayed --> Ready: истёк backoff InFlight --> DLQ: неповторяемая ошибка
невалидная схема, отсутствующая сущность Delayed --> DLQ: попытки исчерпаны DLQ --> Ready: осознанный replay после починки Done --> [*] note right of InFlight Единственное состояние, где два узла могут одновременно считать сообщение своим end note note right of DLQ DLQ без алерта и без процедуры возврата — это /dev/null с иллюзией надёжности end note
Сценарий отказа: лиза короче обработки
Самый коварный режим в очередных брокерах. Видимость сообщения — это обычная лиза со всеми её свойствами (см. координацию): она выдаётся на конечный срок, и её истечение не останавливает того, кто её держал. VisibilityTimeout в SQS — 30 секунд, обработка после деградации внешнего API занимает 45.
очередь «вибрирует»: элементы не уходят,
но и не попадают в DLQ
Диагностические признаки: ApproximateAgeOfOldestMessage растёт, NumberOfMessagesDeleted заметно меньше NumberOfMessagesReceived, в логах — ReceiptHandleIsInvalid или в Kafka-аналоге CommitFailedException. Лечение — три меры вместе: увеличить лизу с запасом относительно p99 обработки; продлевать лизу во время работы (ChangeMessageVisibility в SQS, отдельный heartbeat-поток в RabbitMQ, pause + короткие poll в Kafka); сделать обработчик идемпотентным, потому что первые две меры уменьшают вероятность, но не устраняют её.
Правила, которые экономят инциденты
- Различайте повторяемые и неповторяемые ошибки.
422 Unprocessable, ошибка десериализации, отсутствующая обязательная сущность — не пройдут и с сотой попытки. Их место — DLQ сразу, а не после десяти минут бэкоффа. - Никогда не ретрайте в цикле внутри обработчика, удерживая сообщение. Три попытки по 10 секунд внутри
handle()— это 30 секунд удержания лизы и блокировки партиции. Отдавайте в retry-топик или delayed-очередь и освобождайте поток. - Сохраняйте ключ идемпотентности при переходах между топиками. После DLQ и ручного replay сообщение придёт снова: если ключ сгенерирован заново, дедупликация не сработает.
- Считайте попытки в заголовке сообщения, а не в памяти. Память умирает вместе с процессом; после рестарта счётчик обнулится, и «максимум 5 попыток» станет бесконечностью.
- DLQ обязана иметь алерт и процедуру возврата. Классика жанра: DLQ, в которую за три месяца попало 2,1 млн сообщений, обнаружена при аудите. Алерт на
dlq_depth > 0в течение 15 минут — минимум. Возврат должен быть скриптом, а не ручной операцией в консоли. - Не делайте DLQ входом в бесконечный цикл. Сообщение из DLQ, возвращённое в основной топик без починки причины, вернётся в DLQ. Нужна пометка
replayed_atи отдельный лимит.
Наблюдение за очередями: лаг и его формы
Единственная метрика, которая по-настоящему говорит о здоровье потоковой системы, — лаг, измеренный во времени, а не в сообщениях. «10 000 сообщений» ничего не значит: это 2 секунды при 5000/с и 3 часа при 1/с. Правильная метрика — возраст самого старого необработанного сообщения: ApproximateAgeOfOldestMessage в SQS, kafka_consumergroup_lag_seconds через Burrow или kafka-exporter, разница между now() и timestamp записи в Pulsar.
| Форма кривой лага | Что означает | Что делать |
|---|---|---|
| Пила: рост и падение до нуля | пакетный продюсер, обработка успевает | ничего |
| Ступенька вверх, потом спуск | деплой или рестарт консьюмера | следить, что спуск быстрее интервала деплоев |
| Линейный устойчивый рост | λ > μ, мощности не хватает | масштабировать консьюмеров, потом партиции |
| Плато на высоком уровне | μ ≈ λ, но запаса нет | увеличить мощность до пика |
| Резкое падение до нуля без обработки | ошибочный reset оффсета или retention съел данные | инцидент, проверять auto.offset.reset |
| Лаг ноль и трафика нет | продюсер умер | алерт на отсутствие сообщений обязателен отдельно |
Последняя строка — самая недооценённая. Мониторинг «лаг маленький» показывает зелёный при полностью мёртвом продюсере. Нужен отдельный алерт на rate(messages_in) == 0 дольше ожидаемого интервала.
Автомасштабирование по лагу (KEDA для Kafka/RabbitMQ/SQS) работает, но с жёстким потолком: в Kafka консьюмеров в группе не может быть больше, чем партиций. Масштабирование с 12 до 40 подов при 12 партициях не даст ничего, кроме 28 бездействующих подов и лишнего ребаланса. Увеличение числа партиций — операция, ломающая соответствие «ключ → партиция», планировать её нужно заранее и с запасом.
Как всё это связать в единую картину при разборе инцидента — в наблюдаемости: trace-контекст обязан переживать путешествие через брокер, иначе цепочка разрывается ровно на самом интересном месте.
Почему exactly-once в мессенджинге — обычно миф
Брокер физически не может гарантировать однократное применение эффекта, потому что эффект происходит за его границей. Это тот же аргумент, что и в гарантиях доставки: подтверждение само нуждается в подтверждении, и рекурсия не заканчивается. Но в контексте очередей и потоков стоит перечислить конкретные места, где маркетинговое «exactly-once» перестаёт работать.
- Транзакции Kafka покрывают только Kafka. Атомарность связки «записать в выходные топики + закоммитить входные оффсеты» работает внутри одного кластера. Запись в PostgreSQL, вызов платёжного шлюза, отправка письма в эту транзакцию не входят и войти не могут.
- Мосты между брокерами — всегда at-least-once. MirrorMaker 2, Kafka Connect, RabbitMQ shovel и federation копируют данные между двумя независимыми системами и не могут атомарно закоммитить в обеих. При failover MM2 трансляция оффсетов приблизительная, и дубликаты гарантированы; это прямо написано в документации MirrorMaker.
- Запись «БД → брокер» неатомарна. Классический dual write — сохранили заказ в базе и опубликовали событие — ломается ровно посередине. Единственное корректное решение — transactional outbox или CDC-лог базы; см. распределённые транзакции.
- Fan-out на N подписчиков — это N независимых каналов. Даже если каждый по отдельности «effectively-once», совместной атомарности между ними нет: три из пяти потребителей могут применить событие, два — нет, и это нормальное промежуточное состояние.
- Retention и unclean leader election могут терять подтверждённое.
unclean.leader.election.enable=trueв обмен на доступность разрешает усечение лога — то есть удаление сообщений, за которые продюсер получил ack.
Что делают вместо — пять слоёв, применяемых вместе, а не по выбору:
- Идемпотентный обработчик — база всего: ключ идемпотентности, дедупликация в одной транзакции с эффектом.
- Outbox или CDC на границе «база → брокер», чтобы событие и состояние публиковались атомарно.
- Монотонность и коммутативность там, где можно избавиться от требования порядка вообще.
- Дедупликация с окном на приёмнике — с честным пониманием, что окно конечно и дубликат старше окна пройдёт.
- Сверка (reconciliation) — последний рубеж, который почему-то редко упоминают в статьях про exactly-once, хотя именно он используется в финтехе повсеместно. Периодический процесс сравнивает агрегаты источника и приёмника (число и сумма транзакций за час, контрольные суммы по ключам) и чинит расхождения. Он ловит то, что не поймает ни один брокер: баги в обработчике, ручные вмешательства в базу, потери при усечении лога, сообщения, застрявшие в DLQ.
Формулировка, которую стоит держать в голове при чтении любой документации: exactly-once всегда относится к явно очерченной границе системы. Если в тексте нет предложения, определяющего эту границу, гарантии нет — есть маркетинг. Честный разбор со стороны вендора: Confluent, «Exactly-Once Semantics Are Possible: Here’s How Kafka Does It» — обратите внимание, насколько аккуратно там очерчены границы.
Как выбирать брокер
| Система | Модель | Порядок | Replay | Когда брать |
|---|---|---|---|---|
| Kafka | лог, партиции, ISR | внутри партиции | да, по оффсету | высокий поток, аналитика, event sourcing, много потребителей одних данных |
| RabbitMQ (quorum) | очередь на Raft | FIFO при одном consumer | нет | сложная маршрутизация, приоритеты, тысячи очередей, RPC-подобные сценарии |
| Pulsar | лог поверх BookKeeper | партиция или ключ | да | нужны и лог, и очередные семантики; мультиарендность; геореплика из коробки |
| NATS JetStream | лог, лёгкий | stream | да | низкая задержка, edge-развёртывания, простая эксплуатация, малый след |
| SQS / SNS | очередь managed | FIFO только в FIFO-режиме | нет | AWS-стек, минимум эксплуатации, умеренный поток |
| Redis Streams | лог в памяти | внутри stream | ограниченно | уже есть Redis, допустима потеря при отказе, нужна минимальная задержка |
Три вопроса, которые решают выбор быстрее любой таблицы: нужен ли replay (если да — только лог); нужна ли маршрутизация по содержимому и приоритеты (если да — только очередь); сколько у вас людей на эксплуатацию (если ноль — managed, и точка). Kafka с тремя брокерами, ZooKeeper или KRaft, мониторингом ISR и планированием партиций — это отдельная штатная единица.
Типичные ошибки
- Очередь как «способ ускорить систему». Она не увеличивает μ. Если обработка не успевает — очередь превращает отказ в задержку, и это иногда ровно то, что нужно, но решение должно быть осознанным.
- Неограниченная внутренняя очередь в консьюмере. Разрывает цепь обратного давления и превращает перегрузку в OOM с потерей всего, что было в памяти.
- Порядок, на который никто не рассчитывал. Ключ партиционирования выбран по
user_id, а бизнес-инвариант — на уровнеorder_id. Порядок формально есть, но не тот. - Коммит оффсета до обработки. Превращает at-least-once в at-most-once и даёт тихую потерю при каждом падении.
- Лаг измеряется в сообщениях. График «12 000 сообщений» не отвечает на вопрос «сколько ждёт пользователь».
- Нет алерта на пустой трафик. Мёртвый продюсер выглядит как идеально здоровая система.
- DLQ без владельца. Сообщения копятся, никто не смотрит, через квартал это данные, которые уже нельзя обработать (истёк курс валюты, удалена сущность).
- Ретраи без джиттера и без бюджета. Синхронизированные повторы создают периодические пики и метастабильный отказ, из которого система не выходит сама.
- Изменение числа партиций «на горячую». Соответствие «ключ → партиция» меняется, старые сообщения остаются в старых партициях, порядок нарушен на всю глубину retention.
auto.offset.reset=latestв критичном топике. Однажды это молча пропустит миллионы сообщений, и в логах будет INFO.- Схема сообщений без версионирования. Первое же несовместимое изменение превращается в poison message на всех потребителях сразу.
- Вера в exactly-once из README. Всегда ищите предложение, определяющее границу гарантии.
Мини-итог
- Брокер — это распределённая БД в горячем пути, а не труба. У него есть репликация, лидер, окно долговечности и собственные отказы.
- Два семейства: очередь (состояние на брокере, деструктивное чтение, богатая семантика) и лог (состояние у консьюмера, replay, высокий поток, порядок внутри партиции).
- Глобального порядка не существует. Есть порядок относительно партиции, очереди или ключа. Порядок доставки не равен порядку применения — параллельный консьюмер ломает его без единой ошибки в логах.
- Ключ — одновременно единица порядка и единица параллелизма. Это главный компромисс проектирования топика.
- Backpressure — это сообщение «сбавь темп», а не «выбрось лишнее». Работает, только если ограничено каждое звено; один неограниченный буфер превращает перегрузку в OOM.
- Закон Литтла и
1/(1−ρ)объясняют, почему запас мощности обязателен, а большой буфер вреден. - Повторы устроены по-разному в каждом брокере: считайте попытки в заголовке, продлевайте лизу, разделяйте повторяемые и неповторяемые ошибки, алертите DLQ.
- Exactly-once брокер дать не может. Вместо него — идемпотентность, outbox, коммутативность, дедупликация с окном и сверка как последний рубеж.
Источники
- Kafka Design и Kafka Protocol — первоисточник по модели лога, ISR и семантике доставки.
- Jay Kreps, «The Log: What every software engineer should know about real-time data’s unifying abstraction» — почему лог является универсальной абстракцией.
- Hohpe, Woolf, «Enterprise Integration Patterns» — словарь очередной интеграции: competing consumers, dead letter channel, message ordering.
- Kleppmann, «Designing Data-Intensive Applications», гл. 11 «Stream Processing» — сайт книги.
- RabbitMQ: Quorum Queues и Consumer Prefetch — реальные ручки и их семантика.
- Kyle Kingsbury, Jepsen: RabbitMQ — измеренная потеря данных при разделении сети.
- Uber Engineering, «Building Reliable Reprocessing and Dead Letter Queues with Apache Kafka» — retry-топики в продакшне.
- Reactive Streams Specification — формализация backpressure через
request(n). - Gettys, Nichols, «Bufferbloat: Dark Buffers in the Internet» — почему большой буфер ухудшает систему.
- Amazon SQS: Visibility Timeout и FIFO queues.
- Apache Pulsar: Subscriptions и Key_Shared — порядок по ключу без блокировки партиции.
- Confluent, «Exactly-Once Semantics Are Possible: Here’s How Kafka Does It» — образец аккуратно очерченной границы гарантии.
- KEDA: Kafka scaler и Burrow — мониторинг лага и автомасштабирование.
Что дальше
Мы дошли до точки, где система состоит из брокеров, партиций, ретраев и DLQ — и стала принципиально ненаблюдаемой обычными средствами: стек-трейс обрывается на send(), а причина ошибки находится в другом сервисе двадцать минут спустя. Следующий шаг — научиться видеть распределённую систему целиком.
Наблюдаемость распределённых систем: трассировка, корреляция, поиск причин