Распределённые системы Очереди и потоки: брокеры, порядок сообщений, backpressure, повторы
0%

Очереди и потоки: брокеры, порядок сообщений, backpressure, повторы

Очереди и потоки: брокеры, порядок сообщений, 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» и «данные переживут отказ узла» лежат три независимых решения, и каждое из них по умолчанию настроено в сторону скорости.

  1. Сколько реплик подтвердили. Kafka: acks=0 | 1 | all, причём all означает «все из текущего ISR», который может сжаться до одного узла — поэтому min.insync.replicas обязателен (подробно в репликации). RabbitMQ: publisher confirms плюс quorum queue. SQS: репликация по нескольким AZ выполняется всегда.
  2. Дошло ли до диска. Kafka по умолчанию полагается на page cache и репликацию, а не на fsync каждой записи (flush.messages, flush.ms). Это осознанный выбор: три копии в памяти трёх машин надёжнее одной копии на одном диске — но только если машины действительно независимы (стойка, AZ, питание). Коррелированный отказ ломает эту логику; см. модели отказов.
  3. Подтвердил ли брокер отправителю. Продюсер без 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. Единственный след — жалоба пользователя через неделю.

Лечится это тремя способами, и выбирать нужно осознанно:

  1. Последовательно внутри ключа, параллельно между ключами. Стандартное решение: диспетчеризация по hash(key) % N на N воркеров с собственными очередями.
  2. Монотонность вместо порядка. У события есть версия, применяется только более новая: UPDATE ... WHERE version < :v. Тогда порядок вообще перестаёт быть требованием — подробнее про фенсинг и версии в идемпотентности.
  3. Коммутативные операции. 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 стоит того, чтобы её прочитать целиком (это буквально несколько страниц): она формализует правило, к которому все приходят опытным путём — «производитель не имеет права отправить больше, чем потребитель запросил».

Что видно в логах, когда 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.

Жизненный цикл сообщения с повторами

Сценарий отказа: лиза короче обработки

Самый коварный режим в очередных брокерах. Видимость сообщения — это обычная лиза со всеми её свойствами (см. координацию): она выдаётся на конечный срок, и её истечение не останавливает того, кто её держал. VisibilityTimeout в SQS — 30 секунд, обработка после деградации внешнего API занимает 45.

Диагностические признаки: 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» перестаёт работать.

  1. Транзакции Kafka покрывают только Kafka. Атомарность связки «записать в выходные топики + закоммитить входные оффсеты» работает внутри одного кластера. Запись в PostgreSQL, вызов платёжного шлюза, отправка письма в эту транзакцию не входят и войти не могут.
  2. Мосты между брокерами — всегда at-least-once. MirrorMaker 2, Kafka Connect, RabbitMQ shovel и federation копируют данные между двумя независимыми системами и не могут атомарно закоммитить в обеих. При failover MM2 трансляция оффсетов приблизительная, и дубликаты гарантированы; это прямо написано в документации MirrorMaker.
  3. Запись «БД → брокер» неатомарна. Классический dual write — сохранили заказ в базе и опубликовали событие — ломается ровно посередине. Единственное корректное решение — transactional outbox или CDC-лог базы; см. распределённые транзакции.
  4. Fan-out на N подписчиков — это N независимых каналов. Даже если каждый по отдельности «effectively-once», совместной атомарности между ними нет: три из пяти потребителей могут применить событие, два — нет, и это нормальное промежуточное состояние.
  5. 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 и планированием партиций — это отдельная штатная единица.

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

  1. Очередь как «способ ускорить систему». Она не увеличивает μ. Если обработка не успевает — очередь превращает отказ в задержку, и это иногда ровно то, что нужно, но решение должно быть осознанным.
  2. Неограниченная внутренняя очередь в консьюмере. Разрывает цепь обратного давления и превращает перегрузку в OOM с потерей всего, что было в памяти.
  3. Порядок, на который никто не рассчитывал. Ключ партиционирования выбран по user_id, а бизнес-инвариант — на уровне order_id. Порядок формально есть, но не тот.
  4. Коммит оффсета до обработки. Превращает at-least-once в at-most-once и даёт тихую потерю при каждом падении.
  5. Лаг измеряется в сообщениях. График «12 000 сообщений» не отвечает на вопрос «сколько ждёт пользователь».
  6. Нет алерта на пустой трафик. Мёртвый продюсер выглядит как идеально здоровая система.
  7. DLQ без владельца. Сообщения копятся, никто не смотрит, через квартал это данные, которые уже нельзя обработать (истёк курс валюты, удалена сущность).
  8. Ретраи без джиттера и без бюджета. Синхронизированные повторы создают периодические пики и метастабильный отказ, из которого система не выходит сама.
  9. Изменение числа партиций «на горячую». Соответствие «ключ → партиция» меняется, старые сообщения остаются в старых партициях, порядок нарушен на всю глубину retention.
  10. auto.offset.reset=latest в критичном топике. Однажды это молча пропустит миллионы сообщений, и в логах будет INFO.
  11. Схема сообщений без версионирования. Первое же несовместимое изменение превращается в poison message на всех потребителях сразу.
  12. Вера в exactly-once из README. Всегда ищите предложение, определяющее границу гарантии.

Мини-итог

  • Брокер — это распределённая БД в горячем пути, а не труба. У него есть репликация, лидер, окно долговечности и собственные отказы.
  • Два семейства: очередь (состояние на брокере, деструктивное чтение, богатая семантика) и лог (состояние у консьюмера, replay, высокий поток, порядок внутри партиции).
  • Глобального порядка не существует. Есть порядок относительно партиции, очереди или ключа. Порядок доставки не равен порядку применения — параллельный консьюмер ломает его без единой ошибки в логах.
  • Ключ — одновременно единица порядка и единица параллелизма. Это главный компромисс проектирования топика.
  • Backpressure — это сообщение «сбавь темп», а не «выбрось лишнее». Работает, только если ограничено каждое звено; один неограниченный буфер превращает перегрузку в OOM.
  • Закон Литтла и 1/(1−ρ) объясняют, почему запас мощности обязателен, а большой буфер вреден.
  • Повторы устроены по-разному в каждом брокере: считайте попытки в заголовке, продлевайте лизу, разделяйте повторяемые и неповторяемые ошибки, алертите DLQ.
  • Exactly-once брокер дать не может. Вместо него — идемпотентность, outbox, коммутативность, дедупликация с окном и сверка как последний рубеж.

Источники

Что дальше

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

Наблюдаемость распределённых систем: трассировка, корреляция, поиск причин

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

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

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

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