MAATRIX / Блог / Повторная доставка сообщений списала деньги дважды у 300 клиентов

Повторная доставка сообщений списала деньги дважды у 300 клиентов

MAATRIX

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

Что случилось

Биллинг подписок в сервисе устроен классически: раз в сутки, в 03:00 по московскому времени, cron-джоба находит все подписки, у которых наступил день списания, и кладёт по одному сообщению на каждую в очередь charges.due. Отдельный пул воркеров-консьюмеров разбирает очередь, для каждого сообщения дергает антифрод-проверку, затем вызывает API платёжного шлюза на списание и пишет результат в базу.

В пятницу вечером выкатили изменение: к обработчику списания добавили синхронный вызов сервиса скоринга — раньше решение о блокировке подозрительных платежей принималось асинхронно постфактум, теперь стали блокировать до списания. Сам вызов был простым HTTP-запросом с таймаутом в 3 секунды на стороне клиента. Ревью прошло без замечаний, юнит-тесты зелёные, на стейджинге всё отработало штатно — там нагрузка на антифрод-сервис была на два порядка ниже продовой.

В субботу ночью, когда очередь charges.due разом получила несколько тысяч сообщений (это штатный пик — основная масса подписок продлевается в начале месяца, а конец августа как раз попадает на пиковые даты), антифрод-сервис начал отвечать медленнее. Не критично медленно — p50 остался в норме, но p99 пополз за секунду, а под нагрузкой отдельные запросы занимали 6-8 секунд. Достаточно, чтобы обработка одного сообщения в очереди стала занимать дольше, чем окно, за которое консьюмер обязан подтвердить получение.

Что показали логи и метрики

Первым делом подняли дашборд очереди — в проекте использовался self-hosted NATS JetStream (об устройстве очередей на нём мы уже писали в статье про то, как устроена очередь сообщений). Метрика num_redelivered на консьюмере charges-worker в пятницу ночью и в субботу выглядела аномально: обычно это единицы сообщений в сутки — сеть моргнула, под перезапустился. В ночь инцидента счётчик показывал несколько сотен передоставок за несколько часов.

Дальше подняли логи самого воркера и отфильтровали по order_id тех клиентов, кто написал в саппорт. Картина была одинаковой у всех: два лога charge started и два лога charge succeeded с одним и тем же order_id, с разницей по времени от нескольких секунд до пары минут. При этом message_id у двух записей был разным — то есть это были не два одинаковых лога от гонки внутри одного воркера, а два разных сообщения из очереди, оба с одинаковой полезной нагрузкой.

Отдельно посмотрели метрики самого антифрод-сервиса: график p99 латентности за ночь инцидента показывал явный купол — рост начался примерно через 40 минут после того, как очередь получила основную партию сообщений, и спал ближе к утру, когда нагрузка снизилась. Корреляция с ростом num_redelivered была прямой.

Нужен сервер под эту задачу?

Разверните VPS MAATRIX за пару минут: NVMe, AMD EPYC, root-доступ, локации UK, США, Франция и РФ. Оплата картой РФ и по СБП.

Арендовать сервер

Гипотезы, которые не подтвердились

Прежде чем упереться в реальную причину, отработали три версии, каждая выглядела правдоподобно на старте.

Гонка между воркерами при горизонтальном масштабировании. Под нагрузку автоскейлер поднял дополнительные реплики charges-worker. Первая мысль — два воркера подхватили одно и то же сообщение одновременно. Проверили модель доставки в JetStream: при обычном (не work-queue с ручным конкурентным вычитыванием вне протокола) консьюмере с explicit ack сообщение выдаётся только одному получателю до истечения AckWait, повторная выдача другому воркеру физически невозможна, пока не истёк таймаут. Логи подтвердили: оба дубля по каждому order_id были обработаны одним и тем же инстансом воркера, просто в два захода.

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

Дублирующийся запуск cron-джобы. Проверили, не поднялись ли по ошибке две копии джобы, которая формирует список подписок на списание — например, из-за проблемы с блокировкой при деплое. В логах джобы нашли ровно один запуск за сутки, PID-файл блокировки сработал штатно, второй процесс не стартовал. В очередь charges.due каждое сообщение было положено один раз — продублировалась не публикация, а доставка.

Все три версии закрылись за несколько часов разбора, и стало ясно, что смотреть нужно не на producer и не на платёжный шлюз, а на контракт между консьюмером и брокером.

В чём была реальная причина

У JetStream-консьюмера charges-worker был выставлен AckWait: 5s — время, за которое консьюмер обязан подтвердить обработку сообщения, иначе брокер считает его потерянным и отдаёт заново (столько же секунд, сколько до пятничного деплоя с запасом хватало на всю цепочку: скоринг + запрос к шлюзу + запись в базу). До деплоя средняя обработка укладывалась в 300-600 мс, и даже с запасом на сетевые скачки 5 секунд казались избыточным лимитом.

После деплоя в обработку добавился синхронный вызов антифрод-сервиса. В спокойное время он добавлял 50-150 мс и ничего не менял. Но под пиковой субботней нагрузкой антифрод-сервис начал отвечать за 6-8 секунд на части запросов — то есть дольше, чем весь AckWait. Консьюмер физически ещё выполнял запрос к скорингу, когда JetStream уже решал, что сообщение не подтверждено вовремя, помечал его недоставленным и повторно выдавал — тому же воркеру или другому свободному инстансу, если автоскейлер успел поднять новый под.

Обработчик списания получал вторую копию того же сообщения и не имел никакого способа понять, что это уже обрабатывалось. Кода, который проверял бы «а не списывали ли мы уже по этому order_id за этот расчётный период», в обработчике не было — предполагалось, что раз сообщение одно, значит и вызов один. Итог: первая попытка успевала завершиться и уйти в шлюз на списание уже после того, как брокер решил её передоставить, вторая попытка стартовала параллельно или чуть позже и делала независимый вызов на списание того же счёта.

Ключевая ошибка была не в самом факте передоставки — at-least-once доставка это нормальное и ожидаемое поведение очереди, — а в том, что бизнес-операция со внешними побочными эффектами (списание денег) была реализована как неидемпотентная. Похожая логика разбиралась в статье про очередь задач, которая росла, пока не легло всё: проблема с очередями почти никогда не в самой очереди, а в предположениях, которые обработчик делает о гарантиях доставки.

Как мы это подтвердили

Чтобы не гадать, а доказать причинно-следственную связь, собрали таймлайн по одному конкретному инциденту с дублем. Взяли order_id, по нему нашли оба message_id в логах воркера, а по message_id — записи в JetStream об истории доставки через nats consumer info charges charges-worker --json, где в поле num_redelivered и delivered.consumer_seq видно, что один и тот же stream_seq выдавался консьюмеру дважды с разницей чуть больше пяти секунд — ровно столько, сколько был выставлен AckWait.

Дальше сопоставили это время с логами антифрод-сервиса по трейс-id, который прокидывался через заголовок запроса. Оказалось, что именно на этом запросе латентность составила 7.4 секунды — то есть консьюмер ещё ждал ответ от скоринга, когда брокер уже принял решение о передоставке. Это закрыло вопрос: дело не в баге NATS и не в потерянном ack по сети, а в банальном рассинхроне между временем обработки и таймаутом ожидания подтверждения.

Отдельно подняли MaxDeliver для этого консьюмера — он был выставлен в значение по умолчанию, равное безлимиту в старых конфигурациях, поэтому сообщение могло передоставляться многократно, если бы обработка продолжала подвисать. К счастью, в конкретных инцидентах хватало одной-двух передоставок, но само по себе отсутствие ограничения — отдельная недоработка, которую тоже занесли в список изменений.

Что изменили после

Правки разделили на три уровня: остановить кровотечение, закрыть конкретную причину, снизить вероятность похожего класса ошибок в будущем.

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

Дальше — сделали операцию списания идемпотентной на уровне бизнес-логики, а не только на уровне очереди. Завели таблицу с уникальным ограничением на пару order_id + billing_period:

create table charge_attempts (
    order_id      bigint      not null,
    billing_period date       not null,
    idempotency_key text      not null,
    status        text        not null default 'pending',
    gateway_charge_id text,
    created_at    timestamptz not null default now(),
    unique (order_id, billing_period)
);

Обработчик перед вызовом платёжного шлюза сначала пытается вставить строку с этим ключом:

def process_charge(order_id, billing_period, message_id):
    idempotency_key = f"{order_id}:{billing_period}"
    try:
        insert_charge_attempt(order_id, billing_period, idempotency_key)
    except UniqueViolation:
        # запись уже есть — либо в процессе, либо уже списано
        existing = get_charge_attempt(order_id, billing_period)
        if existing.status == "succeeded":
            log.info("duplicate delivery, charge already succeeded", order_id=order_id)
            return
        # статус pending: используем идемпотентный ключ платёжного шлюза,
        # чтобы даже параллельный запрос не создал второе списание
    result = payment_gateway.charge(
        order_id=order_id,
        idempotency_key=idempotency_key,
    )
    mark_charge_attempt(order_id, billing_period, status="succeeded",
                          gateway_charge_id=result.id)

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

Затем поправили сам таймаут. AckWait подняли до 30 секунд — с запасом на худший случай латентности антифрод-сервиса, а не на средний. Отдельно вынесли вызов скоринга в защищённый вариант: добавили собственный таймаут в 2 секунды на стороне клиента с явной обработкой отказа (при таймауте — деградация к асинхронной постфактум-проверке, а не блокировка всей обработки), плюс circuit breaker, чтобы при деградации антифрод-сервиса воркер не подвешивал каждое сообщение до истечения таймаута. Также выставили MaxDeliver: 5 на консьюмере, чтобы сообщение, которое стабильно не может быть обработано, уходило в dead-letter очередь вместо бесконечных повторных попыток.

Наконец, завели алерт на саму метрику передоставок. Раньше num_redelivered не мониторился вообще — по нему просто не было настроено правило. Теперь при росте передоставок по консьюмеру charges-worker больше чем на 10 в течение 5 минут срабатывает алерт в дежурный канал, и это происходит до того, как накопится статистически значимое число дублей, а не после жалоб в саппорт. Сам разбор инцидента оформили по внутреннему шаблону постмортема — если у вас такого шаблона ещё нет, у нас есть отдельная статья про то, как писать разбор инцидента, с конкретной структурой документа.

Нужен сервер под эту задачу?

Разверните VPS MAATRIX за пару минут: NVMe, AMD EPYC, root-доступ, локации UK, США, Франция и РФ. Оплата картой РФ и по СБП.

Арендовать сервер

Нужны сами нейросети для контента?

Генерируйте изображения, видео и озвучку нейросетями на falapi.io — десятки моделей в одном окне. Оплата картой РФ и по СБП.

Частые вопросы

Почему нельзя было просто сильно увеличить AckWait и не переделывать логику обработчика?

Можно было бы снизить вероятность повтора, но не исключить её полностью — любая очередь с гарантией at-least-once рано или поздно передоставит сообщение: из-за сетевого сбоя, падения пода посреди обработки, рестарта при деплое. Идемпотентность обработчика — это защита от всего класса причин, а не только от конкретного таймаута.

Разве нельзя использовать очередь с гарантией exactly-once и не думать об этом вообще?

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

Как быстро понять, что у вас в проекте есть такой же риск?

Проверьте два места: во-первых, есть ли у любого обработчика очереди, который вызывает платёж, отправку денег, списание баланса или похожую необратимую операцию, проверка на повторную обработку по бизнес-ключу. Во-вторых, посмотрите, мониторится ли у вас метрика передоставок (redelivery count, retry count, DLQ size) — если её нет на дашборде, скорее всего, о повторах в проде вы узнаете только из тикетов.

Что делать, если у сервиса вообще нет метрики передоставок из коробки?

У большинства брокеров она есть, но не всегда включена или не всегда выведена в систему мониторинга по умолчанию. В NATS JetStream это поля num_redelivered и num_ack_pending в выводе consumer info, в RabbitMQ — счётчик redelivered в статистике очереди, в SQS-совместимых системах — ApproximateReceiveCount в атрибутах сообщения. Если брокер крутится на вашем собственном сервере, добавить экспорт такой метрики в Prometheus — обычно вопрос одного конфига экспортера, а не переписывания инфраструктуры.

Стоило ли откатывать деплой с антифрод-проверкой сразу, как только пошли первые жалобы?

Да, и в реальности так и сделали в первые часы — откат сработал как временная остановка кровотечения, пока разбирались в причине. Откат не решает проблему идемпотентности сам по себе (она была скрытой и до этого деплоя), но убирает конкретный триггер, который сделал скрытый риск заметным.

Обсудить статью, задать вопрос или начать новую тему

Есть вопрос по этой статье, идея для обсуждения или просто хотите поделиться опытом? Сообщество MAATRIX ждёт. Для общения, пожалуйста, зарегистрируйтесь в нашем личном кабинете.

Перейти в сообщество →