Главная / Статьи / Таблицы выходящих и входящих сообщений: проектирование webhook’ов, устойчивых к сбоям

Таблицы выходящих и входящих сообщений: проектирование webhook’ов, устойчивых к сбоям

Узнайте, как транзакционная очередь отправлений, идемпотентная очередь входящих сообщений, механизм отложенной обработки с учетом состояния и очереди для некорректных сообщений обеспечивают надежную доставку webhook-сообщений в AWS, Azure и GCP.

2740 слов

Webhooks кажутся самым простым способом интеграции: одна сторона отправляет запрос HTTP POST, а другая его обрабатывает. На практике они несут в себе все риски распределенных систем, поскольку сеть между двумя сервисами может терять запросы, вызывать проблемы из-за истечения времени, доставлять один и тот же пакет данных дважды или перестраивать порядок событий. Если относиться к webhook как к обычному запросу CRUD, в конечном итоге можно потерять уведомления, дважды запустить побочные эффекты и оказаться с двумя системами, которые расходятся во мнениях относительно того, что произошло.

В этом руководстве рассматривается дизайн, способный выдержать такие условия. Вы узнаете, почему примитивный подход терпит неудачу, как транзакционная система отправки сообщений обеспечивает надежность внешних webhook-запросов, как идемпотентная система приема сообщений позволяет безопасно повторять входящие webhook-запросы, как справляться с событиями, поступающими в неправильном порядке, как изолировать данные, которые никогда не могут быть обработаны успешно, и какие управляемые сервисы в AWS, Azure и Google Cloud подходят для каждой части этой архитектуры.

Почему очевидная реализация приводит к потере данных

Рассмотрим SaaS-бэкенд, обрабатывающий значимые изменения состояния, такие как выполнение заказа или активация подписки. Партнер обращается к вашему API для подтверждения действия, и теперь у вашего сервиса появляются две задачи:

  1. Сохранить новое состояние, например, установив статус сущности в Active.
  2. Сообщить сервису нижнего уровня, что сущность готова, отправив ему webhook.

Интуитивный код записывает данные в базу данных, а затем, на следующей строке, отправляет HTTP-запрос. Это и есть двойная запись: две независимые системы обновляются поочередно, при этом между ними нет никакой связи.

Из этого напрямую следуют два варианта сбоя:

  • Процесс прерывается между двумя шагами. База данных считает запись активной, но запрос так и не был отправлен. Ваши данные верны, служба ниже в цепочке не знает об этом, и никто не замечает проблемы до тех пор, пока клиент не подаст жалобу.
  • Запрос отправляется, но транзакция сбивается. Служба ниже в цепочке получает информацию о том, что запись активна, но ваша база данных откатывает операцию и по-прежнему считает её неудачной.

Ни один из вариантов порядка выполнения не решает проблему. Если сначала выполняется HTTP-запрос, возникает вторая ошибка; если же он выполняется в конце — первая. Коренная причина заключается в том, что выполнение коммита базы данных и сетевого запроса нельзя сделать атомарными одновременно, поэтому любая сбой или ошибка между ними приводит к десинхронизации обеих частей.

Надежная отправка вебхуков с использованием транзакционной очереди

Модель очереди устраняет необходимость двойной записи, поскольку HTTP-запрос вообще не выполняется из пути запроса. Вместо этого намерение отправить вебхук преобразуется в данные, которые записываются в той же транзакции базы данных, что и бизнес-изменение. Либо оба операции коммитятся, либо ни одна.

Четыре шага процесса обработки через очередь

  1. Открыть транзакцию. Бизнес-операция запускает обычную транзакцию базы данных.
  • Записывайте состояние и событие вместе. В рамках этой транзакции сервис обновляет таблицу entities (например, устанавливая статус в Active) и вставляет строку в таблицу outbox_events, содержащую точный пакет данных, который должен получить следующий сервис.
  • Фиксируйте транзакцию. После фиксации транзакции база данных надежно сохраняет как новое состояние, так и информацию о необходимости его объявления.
  • Передавайте событие. Отдельный фоновый работник, обычно называемый релеем или публикатором, постоянно ищет строки в таблице outbox, которые еще не были отправлены. Для каждой из них он выполняет запрос HTTP POST, а затем отмечает строку как обработанную.
  • Таблица outbox

    В приведённой ниже таблице хранится по одной строке на каждое ожидаемое уведомление. aggregate_type и aggregate_id указывают, какой бизнес-объект касается событие, event_type определяет, что произошло, payload содержит данные, которые необходимо доставить, а processed_at остаётся пустым до тех пор, пока релеи не подтвердит доставку. Поиск строк с значением null в поле processed_at позволяет получить список задач для релея. Обратите внимание, что встроенные комментарии используют один тире; в PostgreSQL для комментариев требуется два тире (--), поэтому необходимо это исправить перед выполнением запроса.

    CREATE TABLE outbox_events (
    id UUID PRIMARY KEY,
    aggregate_type VARCHAR(50), - e.g., 'Order' or 'User'
    aggregate_id UUID, - e.g., Entity ID
    event_type VARCHAR(100), - e.g., 'order.activated'
    payload JSONB NOT NULL, - The exact webhook payload
    created_at TIMESTAMP DEFAULT NOW(),
    processed_at TIMESTAMP - Null until successfully sent
    );
    

    Что гарантирует очередь отправки, а что нет

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

    Побочным эффектом является то, что доставка происходит хотя бы один раз. Ретранслятор может успешно отправить запрос, а затем выйти из строя до того, как обновит поле processed_at; в таком случае тот же событие будет отправлено снова при следующей обработке. Это допустимо только в том случае, если приемники удаляют дубликаты, что и обеспечивает паттерн почтового ящика с другой стороны. Включение значения id строки из выходного ящика в тело сообщения или заголовок предоставляет приемникам стабильный ключ для идентификации дубликатов. Если вы запускаете несколько экземпляров ретранслятора, убедитесь, что два процесса не смогут одновременно забрать одну и ту же строку; в PostgreSQL для этого часто используется выборка строк с параметром FOR UPDATE SKIP LOCKED.

    Безопасное прием сообщений webhook с использованием идемпотентного почтового ящика

    Теперь рассмотрим ситуацию, когда ваш сервис получает webhook от партнеров или систем вышестоящего уровня.

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

    Паттерн почтового ящика разделяет процесс принятия webhook и его обработку.

    Четыре шага потока в почтовом ящике

    1. Получение и проверка. Как только поступает запрос, проверяется его подпись HMAC, чтобы убедиться, что он действительно пришел от партнера и не был подделан или изменен.
  • Храните необработанный пакет данных. Вставляйте неприкосновенный JSON-тело в таблицу webhook_inbox, используя уникальный идентификатор события партнера в качестве ключа, и защищайте её ограничением уникальности базы данных.
  • Немедленно подтвердите получение. Возвращайте код 200 OK сразу же, до выполнения любой бизнес-логики.
  • Обрабатывайте в фоновом режиме. Рабочий процесс берёт записи из очереди, проверяет, было ли событие уже обработано, пропускает его в случае положительного результата, а в противном случае выполняет бизнес-логику и отмечает запись как обработанную.
  • Таблица inbox

    Здесь каждая строка содержит информацию о том, кто отправил событие (partner_name), идентификатор отправителя (partner_event_id), данные нагрузки, информацию о том, была ли проверена подпись, а также значение status, которое может быть PENDING, PROCESSED или QUARANTINED. Важным элементом является ограничение UNIQUE(partner_name, partner_event_id): именно оно превращает дубликаты в безвредные операции, не имеющие эффекта. Как и в таблице выходящих сообщений, комментарии с одной тирой должны быть заменены на --, чтобы PostgreSQL принял соответствующее заявление.

    CREATE TABLE webhook_inbox (
      id UUID PRIMARY KEY,
      partner_name VARCHAR(50), - e.g., 'Stripe' or 'GitHub'
      partner_event_id VARCHAR(100), - The unique ID from the sender
      payload JSONB NOT NULL,
      signature_verified BOOLEAN,
      status VARCHAR(20), - 'PENDING', 'PROCESSED', 'QUARANTINED'
      received_at TIMESTAMP DEFAULT NOW(),
      processed_at TIMESTAMP,
      UNIQUE(partner_name, partner_event_id) - Prevents duplicate inserts
    );
    

    Почему именно это ограничение выполняет основную работу

    Поскольку обработчик лишь проверяет, вставляет данные и возвращает результат, он реагирует быстро, и у отправителя редко возникают проблемы с таймаутом. Даже если отправитель попытается выполнить операцию повторно, даже десять раз подряд, ограничение на уникальность позволяет успешно выполнить всего одну вставку. Ваш обработчик должен рассматривать возникающую ошибку нарушения уникальности (или результат ON CONFLICT DO NOTHING) как успех и всё равно возвращать код 200 OK; в противном случае отправитель будет продолжать пытаться обработать событие, которое уже имеется. Поскольку существует только одна строка, вспомогательный процесс выполняет свои действия лишь один раз.

    Существуют два момента, которые важно сделать правильно. Во-первых, удаление дубликатов зависит от того, предоставляет ли партнер стабильный идентификатор события; большинство поставщиков webhook включают его, но необходимо подтвердить это для каждой интеграции. Во-вторых, процессор может выйти из строя после выполнения вспомогательных действий, но до того, как будет отмечено, что строка обработана, поэтому там, где это возможно, следует выполнять изменения в бизнес-логике и обновление статуса в рамках одной транзакции, а также делать внешние вспомогательные действия идемпотентными. Чтобы узнать больше о удалении дубликатов запросов с использованием ключей, ознакомьтесь с ключами идемпотентности в POST-конечных точках Node.js.

    Обработка событий, поступающих в неправильном порядке

    Даже при контроле за дубликатами нет гарантии, что события будут поступать в том порядке, в котором они были созданы. Ваш сервис может получить событие entity.completed раньше, чем событие entity.started. Обработчик, который слепо применяет каждое событие, попытается перевести объект из состояния draft непосредственно в состояние completed, что либо повредит его состояние, либо приведёт к ошибке вроде 409 Conflict.

    Проверка каждого перехода с учётом автоматы состояний

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

    На рисунке ниже показано это решение. Если событие завершения поступает, пока объект всё ещё находится в стадии черновика, предварительное условие ещё не выполнено, поэтому функция отмечает событие как отложенное вместо того, чтобы применить его. В комментариях упоминаются два способа обработки отложения: оставить запись в папке с непрочитанными сообщениями и попытаться выполнить операцию позже, либо сохранить прогнозируемое состояние и дождаться отсутствующего события. Событие начала работы для объекта в стадии черновика является допустимой транзицией и применяется. Рассматривайте это как псевдокод: return status: 'DEFERRED'; — недопустимый синтаксис JavaScript, правильно будет return { status: 'DEFERRED' };, а реальная реализация также должна учитывать остальные комбинации событий и состояний.

    function processWebhookEvent(event, currentEntityState) {
        if (event.type === 'entity.completed' && currentEntityState === 'draft') {
            // The 'started' event hasn't arrived yet!
            // We cannot transition from 'draft' directly to 'completed'.
    
            // Option A: Leave it in the inbox and retry in 5 minutes.
            // Option B: Store a "Projected State" and wait for the missing piece.
            return status: 'DEFERRED';
        }
    
        if (event.type === 'entity.started' && currentEntityState === 'draft') {
            return transitionTo('started');
        }
    }
    

    Отложение как цикл самовосстановления

    Возьмем заказ, в котором событие «Отгружено» поступает раньше, чем событие «Оплачено». Мгновенное применение статуса «Отгружено» приведет к тому, что заказ попадет в состояние, которое не допускается вашей моделью. С использованием обработчика, учитывающего состояние заказа, последовательность действий становится следующей:

    1. Поступает событие «Отгружено», оценщик видит, что оплата отсутствует, и событие откладывается.
    2. Поступает событие «Оплачено», оно признано действительным, и заказ обновляется.
    3. Повторно пробуется применить отложенное событие «Отгружено»; теперь условия для его выполнения соблюдены, и оно применяется.

    Задержанные события могут находиться в специальной очереди для повторных попыток, например, в Amazon SQS или очереди на основе Redis, где фоновый работник периодически пытается их обработать. В результате формируется рабочий процесс, который отклоняет недопустимые переходы, но в конечном итоге приходит к правильному состоянию, не удаляя ни одного события. Однако необходимо установить лимит на время, в течение которого событие может оставаться задержанным: если предварительные условия так и не будут выполнены, событие следует в конечном итоге считать неудачей, а не пытаться обработать его бесконечно, что приводит к следующему разделу.

    Изоляция проблемных событий с помощью повторных попыток и очереди для бракованных сообщений

    Некоторые события никогда не будут обработаны, независимо от того, сколько раз вы будете пытаться их выполнить: некорректный формат данных или ссылка на ID, которого нет в вашей базе данных. Это называется «ядовитыми таблетками». Наивный работник будет бесконечно пытаться их обработать, и если очередь обрабатывается по порядку, одно некорректное сообщение может заблокировать все корректные события, находящиеся после него.

    Стандартной мерой защиты является политика ограниченных попыток с увеличивающимися задержками, за которой следует очередь некорректных сообщений (DLQ). Типичный график выглядит следующим образом:

    1. Первая попытка проваливается; ожидание одной минуты.
    2. Вторая попытка проваливается; ожидание пяти минут.
    3. Третья попытка проваливается; ожидание пятнадцати минут.
    4. Четвертая попытка проваливается; событие перемещается в очередь DLQ.

    DLQ может быть таблицей в вашей собственной базе данных или функцией управляемой очереди. Важно то, что происходит дальше: события из DLQ должны отображаться на внутреннем административном интерфейсе и вызывать сигналы высокой приоритетности, поскольку каждое из них представляет собой данные, которые система не смогла обработать. Инженер исследует проблему, устраняет ошибку в сопоставлении или некорректные данные, а затем повторно запускает событие, чтобы оно прошло по обычному пути обработки. Необходимо заранее разработать механизм повторного запуска; без него восстановление после ситуации в DLQ превращается в ручную работу с базой данных под давлением.

    Применение данной концепции в AWS, Azure и Google Cloud

    Папки для отправки и приема сообщений находятся в вашей реляционной базе данных, но сопутствующие компоненты (входные потоки, очереди, процессы обработки, DLQ) хорошо соответствуют управляемым облачным сервисам, что снижает объем операционных задач. Структура остается одинаковой у всех поставщиков; меняются только названия продуктов.

    AWS

    • Входные данные: Amazon API Gateway принимает входящие webhook-сообщения, при этом механизм авторизации Lambda проверяет подпись HMAC до того, как запрос дойдёт до бэкенда.
    • База данных: Amazon Aurora PostgreSQL хранит бизнес-таблицы вместе с таблицами webhook_inbox и outbox_events, что обеспечивает соблюдение транзакционных гарантий.
    • Очереди и очередь ошибок: Стандартная очередь SQS обеспечивает асинхронную обработку, а настроенная очередь ошибок SQS принимает сообщения, когда количество полученных сообщений превышает установленный лимит. Стандартные очереди гарантируют хотя бы однократную обработку сообщений, но не сохраняют их порядок, что ещё раз подчеркивает важность вышеупомянутых проверок идемпотентности и состояния.
  • Рабочие процессы: Функции Lambda, запускаемые SQS для обработки входящих событий. Механизм передачи сообщений работает как запланированная функция Lambda или задача ECS Fargate, которая каждые несколько секунд запрашивает данные из Aurora, отправляет ожидающие события и устанавливает значение processed_at.
  • Azure

    • Входные данные: Azure API Management принимает webhook-сообщения, проверяет их подписи и пересылает запросы на бэкенд.
    • База данных: Azure Database for PostgreSQL Flexible Server хранит состояние приложения, а также таблицы для входящих и исходящих сообщений.
    • Очереди и DLQ: Azure Service Bus координирует передачу сообщений и обеспечивает встроенную систему обработки некорректных сообщений, автоматически откладывая их в отдельную очередь после определенного количества неудачных попыток доставки.
    • Рабочие процессы: Azure Functions с триггерами Service Bus обрабатывают данные из папки входящих сообщений. Механизм отправки сообщений работает в фоновом режиме в Azure Container Apps или как Kubernetes CronJob в случае использования AKS; он запрашивает данные из PostgreSQL на предмет неприсланных событий и передаёт их по протоколу HTTP.

    Google Cloud

    • Входные запросы: Google Cloud API Gateway обрабатывает входящие HTTP-вебхуки и процедуры аутентификации.
    • База данных: Cloud SQL для PostgreSQL хранит реляционные данные, включая обе таблицы.
    • Очереди и DLQ: сервис Pub/Sub асинхронно направляет сообщения. Основная подписка обрабатывает события, а тема для нераспознанных сообщений хранит те сообщения, которые так и не были подтверждены после максимального количества заданных попыток доставки.
    • Рабочие процессы: Сервисы Cloud Run, способные сокращать количество инстансов до нуля между пиковыми нагрузками, получают сообщения через Pub/Sub для обработки входящих данных. Ретранслятор выходящих сообщений может быть задачей Cloud Run или сервисом Cloud Run, запускаемым по расписанию с помощью Cloud Scheduler, который опрашивает Cloud SQL и отправляет ожидающие события.

    Чтобы узнать больше о способах связи между сервисами, таких как OAuth и надежные API-запросы, ознакомьтесь с шестью шаблонами интеграции для надежной связи сервисов Node.js.

    Основные выводы

    • Надежная система webhook — это конвейер обработки событий, а не пара HTTP-конечных точек.
    • Никогда не обновляйте базу данных и не вызывайте удаленный сервис как два независимых шага; записывайте строку в выходную очередь в рамках одной транзакции и пусть ретранслятор доставит ее.
    • Механизм отправки обеспечивает доставку хотя бы один раз, поэтому каждый получатель должен удалять дубликаты.
    • Со стороны получателя необходимо проверить данные, сохранить их с ограничением на уникальность ID события отправителя, немедленно подтвердить получение и выполнять основную работу в фоновом процессе.
    • Проверяйте каждое событие в соответствии с вашей машиной состояний и откладывайте обработку тех событий, у которых отсутствуют предварительные условия, установив лимит на время ожидания.
    • Ограничьте количество попыток повторной отправки, направляйте постоянные сбои в специальную очередь с уведомлениями и сделайте возможность повторной обработки стандартной операцией.
    • Управляемые очереди, такие как SQS, Service Bus и Pub/Sub, обеспечивают функции повторных попыток и обработки сбоев, в то время как таблицы базы данных гарантируют необходимые условия.