Почему очередь доставляет сообщение дважды и почему это не баг
Если вы хоть раз видели в логах одно и то же событие обработанным дважды — заказ отправлен на склад два раза, письмо ушло клиенту второй раз, платёж прошёл повторно, — первая реакция обычно одна: «в очереди баг, сообщения дублируются». Это не баг. Это заявленное поведение почти любого брокера сообщений, который вы используете, и оно называется at least once — «доставлю хотя бы один раз». В этой статье разберём, почему брокеры вообще выбрали такую гарантию вместо честного «ровно один раз», что происходит на уровне протокола в момент дублирования и что с этим знанием делать на практике.
Содержание
- Три гарантии доставки и почему их всего три
- Что происходит в момент, когда сообщение дублируется
- Почему выбрали именно «доставить лишний раз», а не «пропустить»
- Где именно в стеке живёт эта гарантия
- Почему это не решается «просто настройкой брокера понадёжнее»
- Что с этим фактически делать: идемпотентность как единственный надёжный ответ
Три гарантии доставки и почему их всего три
У любого брокера сообщений — RabbitMQ, Kafka, NATS JetStream, SQS — есть ровно три модели гарантии доставки, и это не маркетинговая классификация, а следствие того, что можно физически контролировать в распределённой системе:
- At most once («не более одного раза») — брокер отправил сообщение и сразу забыл о нём. Если потребитель упал до обработки — сообщение потеряно навсегда. Никаких повторов.
- At least once («хотя бы один раз») — брокер хранит сообщение, пока не получит явное подтверждение (ack) от потребителя. Нет подтверждения — сообщение доставляется снова. Потерь нет, но возможны дубликаты.
- Exactly once («ровно один раз») — сообщение доставляется и обрабатывается один и только один раз, без потерь и без дублей.
На первый взгляд exactly once выглядит идеальным вариантом, и разработчики систем закономерно спрашивают: почему бы брокеру не делать сразу так? Ответ лежит в природе распределённых систем: сеть между потребителем и брокером ненадёжна, и вопрос «а точно ли сообщение было обработано» в общем случае неразрешим, если сеть может как угодно задерживать и терять пакеты. Это не инженерная лень — это теоретическое ограничение, знакомое как проблема двух генералов и позже сформулированное более строго в контексте распределённого консенсуса. Kafka и некоторые другие системы предлагают exactly once semantics для отдельных сценариев (например, обработка внутри одного кластера Kafka от источника до приёмника через transactional producer), но это узкая, дорогая по ресурсам гарантия для конкретного паттерна использования — а не универсальное свойство доставки сообщений между независимыми системами. В общем случае, когда потребитель — это ваш процесс, а брокер — отдельная система за сетью, get exactly once честно не получится, и все крупные брокеры сообщений это открыто признают в своей документации.
Что происходит в момент, когда сообщение дублируется
Разберём механику по шагам — это последовательность, которая объясняет буквально каждый случай дублирования, который вы видели в логах.
- Потребитель забирает сообщение из очереди (или из партиции топика). Брокер помечает его как «выдано, ожидает подтверждения» и запускает таймер ожидания (visibility timeout в SQS, unacked message в RabbitMQ, отсутствие коммита offset в Kafka).
- Потребитель обрабатывает сообщение — например, списывает деньги со счёта или создаёт запись в БД. Обработка завершена успешно.
- Потребитель отправляет брокеру подтверждение (ack, commit offset, delete message).
- Подтверждение не доходит до брокера. Причины могут быть разные: оборвалось TCP-соединение, между потребителем и брокером произошёл сетевой сбой, сам процесс потребителя упал сразу после обработки, но до отправки ack, либо истёк таймаут ожидания подтверждения, потому что обработка заняла дольше, чем брокер готов был ждать.
- Брокер, не получив ack вовремя, оказывается в ситуации фундаментальной неопределённости: он не может отличить «потребитель обработал сообщение, но ack потерялся в сети» от «потребитель вообще не успел обработать сообщение». С точки зрения брокера это два неразличимых состояния — у него просто нет сообщения-подтверждения, и точка.
- У брокера в этот момент есть ровно два варианта поведения. Первый — считать, что раз ack не пришёл, сообщение, скорее всего, не обработано, и не повторять доставку. Второй — считать, что раз нет уверенности в обработке, надёжнее доставить сообщение ещё раз. Разработчики протоколов брокеров почти единогласно выбрали второй вариант, потому что цена ошибки в нём меньше.
Ключевой момент здесь — шаг 5. Дублирование происходит не потому, что брокер «глючит» и путает сообщения. Оно происходит потому, что у брокера физически нет способа узнать правду о состоянии потребителя в момент истечения таймаута, а решение нужно принять. Это в чистом виде проблема консенсуса при ненадёжной сети, и брокер выбирает консервативную сторону.
Нужен сервер под эту задачу?
Разверните VPS MAATRIX за пару минут: NVMe, AMD EPYC, root-доступ, локации UK, США, Франция и РФ. Оплата картой РФ и по СБП.
Арендовать серверПочему выбрали именно «доставить лишний раз», а не «пропустить»
Представьте, что вы проектируете протокол доставки и стоите перед этим выбором. Разберём цену ошибки в обе стороны на конкретном примере — обработка платежа через очередь.
Если брокер выбирает не повторять при сомнении (ближе к at most once): в сценарии из шага 4 сообщение считается доставленным и удаляется из очереди, даже если потребитель на самом деле не успел его обработать (упал до записи в БД). Итог — платёж просто пропал. Клиент нажал «оплатить», деньги списались (или не списались) на стороне платёжного шлюза, но система заказов никогда не узнает об этом событии. Восстановить такую потерю практически невозможно — сообщения больше нет нигде, и единственный способ его найти — сверка со внешней системой (если она вообще ведётся) постфактум.
Если брокер выбирает повторить при сомнении (at least once): в том же сценарии сообщение будет доставлено снова. Если потребитель на самом деле уже обработал его в первый раз — вы получаете дубль: платёж может обработаться дважды. Это неприятно, но это ошибка, которую видно, которую можно поймать и исправить на уровне обработчика (о чём ниже) — и вы точно знаете, что событие не потеряно.
Асимметрия здесь принципиальная: потерянное сообщение исчезает бесследно и часто необнаружимо, а дублированное сообщение видно и обрабатываемо. Для подавляющего большинства бизнес-сценариев — платежи, заказы, уведомления, синхронизация данных — потеря события дороже, чем повторная обработка. Поэтому индустриальный консенсус (RabbitMQ, Kafka, SQS, Google Pub/Sub, Azure Service Bus) сошёлся именно на at least once как гарантии по умолчанию. Это не случайность и не недоработка — это результат осознанного взвешивания рисков, сделанного авторами протоколов задолго до того, как вы установили свой брокер.
Здесь же стоит развеять смежное заблуждение: at most once иногда всё-таки применяют осознанно, но только там, где потеря сообщения дешевле, чем задержка или лишняя нагрузка от повтора — например, в метриках телеметрии, где одно потерянное измерение из миллиона не имеет значения, а повторная доставка миллионов метрик создаёт лишнюю нагрузку на систему сбора. Это редкий случай, и выбирать его нужно осознанно для конкретного потока данных, а не как гарантию по умолчанию для всей системы.
Где именно в стеке живёт эта гарантия
Важно понимать, что at least once — это гарантия уровня «брокер — потребитель», и она не защищает автоматически весь путь сообщения от начала до конца. Разберём, где именно она действует в каждой из популярных систем.
| Система | Что подтверждает потребитель | Что произойдёт при потере ack |
|---|---|---|
| RabbitMQ | basic.ack на конкретное сообщение | Сообщение остаётся unacked, после разрыва соединения или requeue возвращается в очередь |
| Kafka | commit offset (сам номер позиции в партиции, а не отдельное сообщение) | Начиная со следующего чтения консьюмер снова получит все сообщения от последнего закоммиченного offset |
| Amazon SQS | DeleteMessage с receipt handle | Сообщение снова становится видимым после истечения visibility timeout |
| NATS JetStream | Ack() на сообщение с явным подтверждением | Сообщение повторно доставляется после AckWait |
Обратите внимание на Kafka отдельно: там нет подтверждения на уровне отдельного сообщения — есть коммит offset, то есть «я обработал всё вплоть до этой позиции». Если вы обработали пять сообщений подряд и закоммитили offset только после пятого, а процесс упал на третьем — при перезапуске вы получите заново все пять, включая уже обработанные первые два. Это тот же принцип at least once, только гранулярность подтверждения другая — не по одному сообщению, а по позиции в логе. Подробнее про устройство продюсера, консьюмера и партиций можно почитать в статье о том, как устроена очередь сообщений.
Отдельный нюанс: гарантия at least once описывает только путь от брокера до потребителя. Она ничего не говорит о том, что происходит до брокера (продюсер тоже может отправить сообщение дважды, если не получил подтверждение о записи) и что происходит после успешной обработки потребителем (если сам бизнес-эффект — например, запись в БД — не атомарен с отправкой ack, можно точно так же потерять или задублировать эффект независимо от гарантий брокера).
Почему это не решается «просто настройкой брокера понадёжнее»
Возникает соблазн подумать: раз проблема в потере ack, может, есть режим брокера, где ack точно не теряется? Здесь важно понимать границу возможного.
Можно снизить вероятность потери ack — более надёжное соединение, разумные таймауты, retry на уровне клиентской библиотеки. Это уменьшает частоту лишних повторов. Но исключить сценарий нельзя в принципе: асинхронная сеть допускает произвольные задержки, а гарантированно отличить «ack задержался» от «ack потерян навсегда» без бесконечного ожидания невозможно. Любой конечный таймаут — компромисс между скоростью реакции и точностью решения.
Увеличение таймаута ack не устраняет проблему, а сдвигает баланс между двумя издержками: слишком короткий таймаут — брокер рано решает, что ack потерян, и шлёт дубликат, хотя потребитель просто дольше обрабатывал сообщение; слишком длинный — при реальном падении потребителя сообщение долго остаётся «зависшим». Настройка таймаутов (visibility timeout в SQS, AckWait в JetStream, session.timeout.ms в Kafka) снижает частоту ложных повторов, но не убирает саму возможность дублирования — убрать её значило бы решить нерешаемую в общем случае задачу консенсуса при ненадёжной сети.
Ещё один источник дублей: rebalance партиций в Kafka или переподключение consumer group. При ребалансировке офсеты, не закоммиченные на этот момент, приведут к повторной обработке этих сообщений новым владельцем партиции. Механизм тот же — брокер не уверен, что сообщение обработано, раз подтверждения не было, — но триггер другой: не сетевой сбой, а плановое перераспределение нагрузки.
Что с этим фактически делать: идемпотентность как единственный надёжный ответ
Раз дублирование убрать нельзя, инженерный ответ смещается с «как избежать дублей» на «как сделать так, чтобы дубль не приводил к повторному эффекту». Это свойство называется идемпотентностью обработчика: повторный вызов с теми же данными не должен менять результат по сравнению с однократным вызовом.
Практические способы добиться этого без переписывания всей бизнес-логики:
- Уникальный идентификатор события плюс таблица обработанных ID. Продюсер кладёт в сообщение уникальный
message_id(UUID, детерминированно сгенерированный по содержимому события). Потребитель перед обработкой проверяет по таблице или ключу в Redis, не обрабатывал ли он уже этот ID, и если да — просто подтверждает получение (ack) без повторного применения эффекта.
-- пример таблицы дедупликации в PostgreSQL
CREATE TABLE processed_events (
event_id UUID PRIMARY KEY,
processed_at TIMESTAMPTZ NOT NULL DEFAULT now()
);
-- обработчик перед применением эффекта
INSERT INTO processed_events (event_id) VALUES ($1)
ON CONFLICT (event_id) DO NOTHING
RETURNING event_id;
-- если строка не вернулась — событие уже обработано, эффект не применяем
- UPSERT вместо INSERT там, где это применимо. Если эффект обработки — запись состояния (а не накопление, вроде списания баланса), замена
INSERTнаINSERT ... ON CONFLICT DO UPDATEделает повторное применение безопасным автоматически: второй раз записывается то же самое значение.
- Идемпотентные ключи на уровне внешних API. Многие платёжные системы (это общепринятая практика в индустрии, а не особенность конкретного провайдера) поддерживают
Idempotency-Keyв заголовке запроса — вы передаёте туда тот жеmessage_id, и повторный запрос с этим ключом возвращает результат первого вызова, не выполняя списание повторно.
- Атомарность эффекта и подтверждения обработки. Если возможно, применение эффекта (например, вставка записи в БД) и запись факта обработки должны происходить в одной транзакции — тогда даже при падении процесса между этими шагами не возникает рассинхронизации, и повторная обработка сообщения безопасно определяется по той же таблице.
- Естественная идемпотентность операции там, где это возможно проектно. Операция «установить баланс равным 500» идемпотентна сама по себе; операция «прибавить 500 к балансу» — нет. Там, где бизнес-логика позволяет, лучше проектировать события как явные состояния, а не как дельты.
Дедупликация по message_id надёжна только в пределах окна, за которое вы храните обработанные ID — если хранить их вечно, таблица растёт бесконечно. На практике ставят TTL на запись (например, 7–30 дней, в зависимости от максимально возможной задержки повторной доставки) либо используют bloom-фильтр для компактной приближённой проверки.
Нужен сервер под эту задачу?
Разверните VPS MAATRIX за пару минут: NVMe, AMD EPYC, root-доступ, локации UK, США, Франция и РФ. Оплата картой РФ и по СБП.
Арендовать серверНужны сами нейросети для контента?
Генерируйте изображения, видео и озвучку нейросетями на falapi.io — десятки моделей в одном окне. Оплата картой РФ и по СБП.
Частые вопросы
Если у меня Kafka с exactly once semantics (EOS), можно ли не думать об идемпотентности?
Только если весь путь данных — от источника до приёмника — целиком лежит внутри Kafka и обрабатывается через transactional producer с isolation.level=read_committed. Как только в цепочке появляется внешняя система (ваша БД, HTTP-вызов, другой брокер), гарантия EOS Kafka на неё не распространяется, и идемпотентность обработчика снова становится вашей ответственностью.
Можно ли просто увеличить таймаут ack, чтобы дублей стало меньше?
Да, это снижает частоту ложных повторов, вызванных тем, что обработка не успела уложиться в таймаут. Но это не убирает саму возможность дублирования при реальном сетевом сбое, а слишком большой таймаут увеличивает время, за которое зависшее сообщение переходит другому обработчику при настоящем падении процесса.
RabbitMQ поддерживает at most once, если выключить подтверждения (auto-ack)?
Да, при noAck: true брокер удаляет сообщение из очереди сразу при отправке, не дожидаясь подтверждения. Это ускоряет обработку, но означает, что при падении потребителя между получением и обработкой сообщение теряется безвозвратно — используйте этот режим только там, где потеря отдельных сообщений допустима.
Как отличить дубль, вызванный потерей ack, от дубля, который продюсер отправил сам по ошибке?
С точки зрения потребителя — никак, оба выглядят как повторное сообщение с тем же (или разным) message_id. Если продюсер тоже должен быть идемпотентным (например, ретраит отправку при таймауте записи), убедитесь, что он передаёт детерминированный message_id, привязанный к содержимому события, а не генерирует новый UUID при каждой попытке — иначе дедупликация на стороне потребителя не сработает.
Есть ли брокеры, которые гарантируют exactly once «из коробки» для любого сценария?
Нет ни одного брокера, который решал бы эту задачу универсально для произвольной внешней системы на другом конце — это следствие теоретической невозможности консенсуса при ненадёжной сети без дополнительных допущений, а не ограничение конкретного продукта. Некоторые системы (Kafka Streams, Google Cloud Pub/Sub с exactly-once delivery в определённых конфигурациях) предлагают близкие к exactly once гарантии для узкого класса сценариев внутри своей экосистемы, но всегда с оговорками по границам применимости.
Обсудить статью, задать вопрос или начать новую тему
Есть вопрос по этой статье, идея для обсуждения или просто хотите поделиться опытом? Сообщество MAATRIX ждёт. Для общения, пожалуйста, зарегистрируйтесь в нашем личном кабинете.
Перейти в сообщество →