Создание системы доставки веб-хуков: надежные постбэки в 2026 году

Практическое руководство 2026 года по обеспечению надежной доставки веб-хуков и партнерских постбэков с использованием очередей, неравномерных повторных попыток, автоматических выключателей, идемпотентности, подписей и телеметрии задержки для каждой конечной точки.
Система доставки веб-хуков — Создание системы доставки веб-хуков: надежные обратные запросы в 2026 году

Последнее обновление 24 июня 2026 г. Цезарь Фиксон

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

В этой реализации намеренно используется небольшая модель выполнения: рабочие процессы на Python и PostgreSQL для очереди, истории доставок и проверки работоспособности конечных точек. Это практичная отправная точка для B2B SaaS, партнерских программ в сфере iGaming и любого продукта, которому необходимо обрабатывать исходящие события, не рассматривая ни один HTTP-таймаут как потерянную конверсию.

Что должно гарантировать договор поставки

Управление Почему это важно Оперативная проверка
Долговечная запись событий События сохраняются после перезапуска рабочих процессов. Каждая принятая посылка имеет стабильный идентификатор и статус.
Подписанный запрос Получатели могут проверить отправителя и обнаружить попытки несанкционированного доступа к данным. Используйте подпись с отметкой времени и отклоняйте устаревшие запросы.
Защита от дублирования Повторные попытки могут привести к дублированию преобразований или обновлений. Отправьте идентификатор события и сделайте получателя идемпотентным.
Ограниченное количество повторных попыток Временные сбои позволяют восстановить работоспособность устройства без перегрузки поврежденного конечного устройства. Сбавьте темп, используя метод дрожания; остановитесь после достижения установленного лимита попыток.
Конечная телеметрия Одна лишь глубина очереди скрывает медленную работу или сбои в работе партнеров. Отслеживайте процент успешных запросов, самое старое ожидающее событие и процентные показатели задержки для каждой конечной точки.

Вопросы безопасности и контроля дубликатов рассматриваются до настройки количества повторных попыток.

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

В рамках партнерского или iGaming-процесса обработки платежей, идентификаторы кликов, идентификаторы конверсий и поля, относящиеся к выплатам, не следует заносить в журналы, если только явно не указаны правила контроля доступа и хранения данных. Приведенная ниже ссылка на Scaleo актуальна только в качестве примера задержки обработки платежей в партнерской программе; она не заменяет документирование договора о предоставлении услуг между вашими собственными сервисами.

Полезные примеры реализации: проверка подписи веб-перехватчика Stripe, PostgreSQL SELECT and SKIP LOCKED и Рекомендации AWS по экспоненциальной задержке и дрожанию.

Модель данных

Всё начинается с двух таблиц: одна для очереди доставки, а другая для измерения задержки.

SQL

-- Pending and in-flight webhook deliveries
CREATE TABLE webhook_queue (
    id              SERIAL PRIMARY KEY,
    endpoint_id     INTEGER NOT NULL,
    endpoint_url    TEXT NOT NULL,
    payload         JSONB NOT NULL,
    
    -- Delivery state
    status          VARCHAR(20) NOT NULL DEFAULT 'pending',
        -- pending, in_flight, delivered, failed, dead_letter
    attempt_count   INTEGER NOT NULL DEFAULT 0,
    max_attempts    INTEGER NOT NULL DEFAULT 5,
    
    -- Timing
    created_at      TIMESTAMP NOT NULL DEFAULT NOW(),
    next_attempt_at TIMESTAMP NOT NULL DEFAULT NOW(),
    delivered_at    TIMESTAMP,
    last_error      TEXT,
    
    -- Response tracking
    last_status_code INTEGER,
    last_response_ms INTEGER,  -- response time in milliseconds
    
    INDEX idx_queue_next (status, next_attempt_at)
        WHERE status IN ('pending', 'failed')
);

-- Per-endpoint latency measurements (ring buffer)
CREATE TABLE endpoint_latency (
    id              SERIAL PRIMARY KEY,
    endpoint_id     INTEGER NOT NULL,
    response_ms     INTEGER NOT NULL,
    status_code     INTEGER,
    measured_at     TIMESTAMP NOT NULL DEFAULT NOW(),
    
    INDEX idx_latency_endpoint (endpoint_id, measured_at)
);

-- Endpoint health state (circuit breaker)
CREATE TABLE endpoint_health (
    endpoint_id         INTEGER PRIMARY KEY,
    endpoint_url        TEXT NOT NULL,
    state               VARCHAR(20) NOT NULL DEFAULT 'closed',
        -- closed (healthy), open (broken), half_open (testing)
    consecutive_failures INTEGER NOT NULL DEFAULT 0,
    failure_threshold   INTEGER NOT NULL DEFAULT 10,
    last_failure_at     TIMESTAMP,
    last_success_at     TIMESTAMP,
    opened_at           TIMESTAMP,  -- when circuit opened
    cooldown_seconds    INTEGER NOT NULL DEFAULT 300,  -- 5 min before half_open
    
    -- Latency stats (updated periodically)
    p50_ms              INTEGER,
    p95_ms              INTEGER,
    p99_ms              INTEGER,
    sample_count        INTEGER NOT NULL DEFAULT 0
);

Следует отметить три момента относительно этой схемы.

Во-первых, webhook_queue в таблице используется next_attempt_at Вместо отдельного механизма планирования используется столбец. Рабочий процесс опрашивает строки, где status IN ('pending', 'failed') AND next_attempt_at <= NOW()Это упрощенная очередь с задержкой, которая хорошо работает примерно до 10 000 доставок в минуту. При более высоких скоростях лучше использовать полноценный брокер сообщений.

Во-вторых, endpoint_latency Таблица действует как кольцевой буфер. Я периодически удаляю строки старше 24 часов. Процентили задержки в endpoint_health Они рассчитываются на основе этого скользящего окна — они отражают недавнее поведение, а не исторические средние значения.

В-третьих, endpoint_health Эта таблица реализует конечный автомат автоматического выключателя. Подробнее об этом ниже.

Курьер

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

питон

import requests
import time
import psycopg2
from psycopg2.extras import RealDictCursor
from datetime import datetime, timedelta

DB_DSN = "postgresql://user:pass@localhost/webhooks"

def get_connection():
    return psycopg2.connect(DB_DSN)

def deliver_webhooks(batch_size=50):
    """
    Fetch pending webhooks and attempt delivery.
    Uses SELECT FOR UPDATE SKIP LOCKED for safe concurrent workers.
    """
    conn = get_connection()
    cur = conn.cursor(cursor_factory=RealDictCursor)
    
    try:
        cur.execute("""
            SELECT id, endpoint_id, endpoint_url, payload, 
                   attempt_count, max_attempts
            FROM webhook_queue
            WHERE status IN ('pending', 'failed')
              AND next_attempt_at <= NOW()
            ORDER BY next_attempt_at ASC
            LIMIT %s
            FOR UPDATE SKIP LOCKED
        """, (batch_size,))
        
        rows = cur.fetchall()
        
        for row in rows:
            # Check circuit breaker before attempting
            if is_circuit_open(cur, row['endpoint_id']):
                # Don't attempt delivery — reschedule for after cooldown
                reschedule_for_cooldown(cur, row['id'], row['endpoint_id'])
                continue
            
            # Attempt delivery and measure latency
            result = attempt_delivery(
                row['endpoint_url'], 
                row['payload']
            )
            
            # Record latency measurement regardless of success/failure
            record_latency(
                cur, 
                row['endpoint_id'], 
                result['response_ms'], 
                result['status_code']
            )
            
            if result['success']:
                mark_delivered(cur, row['id'], result)
                record_success(cur, row['endpoint_id'])
            else:
                handle_failure(
                    cur, row['id'], row['endpoint_id'],
                    row['attempt_count'], row['max_attempts'],
                    result
                )
        
        conn.commit()
    
    except Exception as e:
        conn.rollback()
        raise
    finally:
        cur.close()
        conn.close()


def attempt_delivery(url, payload):
    """
    Fire the webhook and measure response time.
    Returns dict with success, status_code, response_ms, error.
    """
    start = time.monotonic()
    
    try:
        response = requests.post(
            url,
            json=payload,
            timeout=15,          # 15 second hard timeout
            headers={
                'Content-Type': 'application/json',
                'User-Agent': 'WebhookDelivery/1.0',
                'X-Delivery-Timestamp': str(int(time.time()))
            }
        )
        
        elapsed_ms = int((time.monotonic() - start) * 1000)
        
        return {
            'success': 200 <= response.status_code < 300,
            'status_code': response.status_code,
            'response_ms': elapsed_ms,
            'error': None if response.ok else f"HTTP {response.status_code}"
        }
    
    except requests.Timeout:
        elapsed_ms = int((time.monotonic() - start) * 1000)
        return {
            'success': False,
            'status_code': None,
            'response_ms': elapsed_ms,
            'error': 'timeout_15s'
        }
    
    except requests.ConnectionError as e:
        elapsed_ms = int((time.monotonic() - start) * 1000)
        return {
            'success': False,
            'status_code': None,
            'response_ms': elapsed_ms,
            'error': f'connection_error: {str(e)[:200]}'
        }

SELECT FOR UPDATE SKIP LOCKED Этот пункт имеет решающее значение для запуска нескольких экземпляров рабочих процессов. Без него... SKIP LOCKEDВ этом случае два рабочих процесса будут блокироваться на одной и той же строке. При этом каждый рабочий процесс будет обрабатывать разные пакеты ожидающих веб-хуков. Это обеспечивает горизонтальное масштабирование за счет простого запуска большего количества рабочих процессов.

time.monotonic() звонить вместо time.time() Это сделано намеренно. time.time() Может произойти скачок назад во время корректировки NTP. time.monotonic() Движение никогда не возвращается назад, что важно при измерении задержки менее секунды.

Логика повторных попыток с экспоненциальной задержкой и дрожанием сигнала

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

питон

import random

def calculate_next_attempt(attempt_count, base_delay=30, max_delay=3600):
    """
    Exponential backoff with full jitter.
    
    Attempt 1: 0-30s
    Attempt 2: 0-60s  
    Attempt 3: 0-120s
    Attempt 4: 0-240s
    Attempt 5: 0-480s (capped at max_delay)
    
    Full jitter prevents thundering herd when an endpoint
    recovers and hundreds of retries fire simultaneously.
    """
    exponential_delay = base_delay * (2 ** attempt_count)
    capped_delay = min(exponential_delay, max_delay)
    jittered_delay = random.uniform(0, capped_delay)
    
    return datetime.utcnow() + timedelta(seconds=jittered_delay)


def handle_failure(cur, webhook_id, endpoint_id, 
                   attempt_count, max_attempts, result):
    """
    Handle a failed delivery attempt.
    Either retry with backoff or move to dead letter queue.
    """
    new_attempt_count = attempt_count + 1
    
    if new_attempt_count >= max_attempts:
        # Exhausted retries — dead letter
        cur.execute("""
            UPDATE webhook_queue 
            SET status = 'dead_letter',
                attempt_count = %s,
                last_error = %s,
                last_status_code = %s,
                last_response_ms = %s
            WHERE id = %s
        """, (
            new_attempt_count, result['error'],
            result['status_code'], result['response_ms'],
            webhook_id
        ))
    else:
        # Schedule retry with backoff
        next_attempt = calculate_next_attempt(new_attempt_count)
        cur.execute("""
            UPDATE webhook_queue
            SET status = 'failed',
                attempt_count = %s,
                next_attempt_at = %s,
                last_error = %s,
                last_status_code = %s,
                last_response_ms = %s
            WHERE id = %s
        """, (
            new_attempt_count, next_attempt,
            result['error'], result['status_code'],
            result['response_ms'], webhook_id
        ))
    
    # Update circuit breaker
    record_failure(cur, endpoint_id)

Почему используется полный джиттер, а не декоррелированный или равный джиттер? AWS опубликовала исчерпывающий анализ по этому вопросу. Полный джиттер (случайное изменение значения от 0 до экспоненциального предела) обеспечивает наименьшее общее время выполнения для всех клиентов. Равный джиттер (случайное изменение значения от половины предела до полного предела) более консервативен, но медленнее обрабатывает очередь повторных попыток. Для доставки веб-хуков, где у вас много независимых конечных точек, полный джиттер — правильный выбор, потому что повторные попытки для каждой конечной точки независимы — между ними нет координации.

Автоматический выключатель: Прекратите стучать по неисправным конечным точкам

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

питон

def is_circuit_open(cur, endpoint_id):
    """
    Check if the circuit breaker is open (endpoint is broken).
    If open and cooldown has passed, transition to half_open.
    """
    cur.execute("""
        SELECT state, opened_at, cooldown_seconds
        FROM endpoint_health
        WHERE endpoint_id = %s
    """, (endpoint_id,))
    
    row = cur.fetchone()
    if not row:
        return False  # no health record = assume healthy
    
    if row['state'] == 'closed':
        return False
    
    if row['state'] == 'open':
        # Check if cooldown period has elapsed
        if row['opened_at'] and row['cooldown_seconds']:
            elapsed = (datetime.utcnow() - row['opened_at']).total_seconds()
            if elapsed >= row['cooldown_seconds']:
                # Transition to half_open — allow one probe
                cur.execute("""
                    UPDATE endpoint_health
                    SET state = 'half_open'
                    WHERE endpoint_id = %s
                """, (endpoint_id,))
                return False  # allow the probe delivery
        return True  # still in cooldown
    
    if row['state'] == 'half_open':
        return False  # allow probe delivery
    
    return False


def record_failure(cur, endpoint_id):
    """
    Record a delivery failure. Open circuit if threshold reached.
    """
    cur.execute("""
        UPDATE endpoint_health
        SET consecutive_failures = consecutive_failures + 1,
            last_failure_at = NOW()
        WHERE endpoint_id = %s
        RETURNING consecutive_failures, failure_threshold, state
    """, (endpoint_id,))
    
    row = cur.fetchone()
    if not row:
        # Create health record on first failure
        cur.execute("""
            INSERT INTO endpoint_health (endpoint_id, endpoint_url, 
                consecutive_failures, last_failure_at)
            VALUES (%s, '', 1, NOW())
            ON CONFLICT (endpoint_id) DO UPDATE
            SET consecutive_failures = endpoint_health.consecutive_failures + 1,
                last_failure_at = NOW()
        """, (endpoint_id,))
        return
    
    if row['state'] == 'half_open':
        # Probe failed — re-open circuit with longer cooldown
        cur.execute("""
            UPDATE endpoint_health
            SET state = 'open',
                opened_at = NOW(),
                cooldown_seconds = LEAST(cooldown_seconds * 2, 3600)
            WHERE endpoint_id = %s
        """, (endpoint_id,))
    
    elif (row['state'] == 'closed' and 
          row['consecutive_failures'] >= row['failure_threshold']):
        # Threshold reached — open circuit
        cur.execute("""
            UPDATE endpoint_health
            SET state = 'open',
                opened_at = NOW(),
                cooldown_seconds = 300  -- reset to 5 minutes
            WHERE endpoint_id = %s
        """, (endpoint_id,))


def record_success(cur, endpoint_id):
    """
    Record a delivery success. Close circuit if half_open.
    """
    cur.execute("""
        UPDATE endpoint_health
        SET consecutive_failures = 0,
            last_success_at = NOW(),
            state = 'closed'
        WHERE endpoint_id = %s
    """, (endpoint_id,))

В большинстве реализаций упускается из виду деталь, касающаяся перехода от состояния «половина открыта» к состоянию «открыта» с удвоенным временем ожидания. Если конечная точка выходит из строя во время проверки (состояние «половина открыта»), повторная попытка через 5 минут нежелательна. Конечная точка по-прежнему неисправна. Удвойте время ожидания до 10 минут, затем до 20, и доведите до 1 часа. Это предотвратит превращение автоматического выключателя в периодический механизм принудительного отключения.

Отслеживание процентилей задержки

Средние значения вводят в заблуждение. Конечная точка со средним временем отклика 200 мс может отвечать за 50 мс в 95% случаев и за 3,000 мс в остальных 5%. В среднем всё выглядит нормально. Показатель P95 выявляет проблему, которая затрагивает 1 из 20 доставок.

питон

def record_latency(cur, endpoint_id, response_ms, status_code):
    """
    Record a latency measurement and update percentile stats.
    """
    cur.execute("""
        INSERT INTO endpoint_latency 
            (endpoint_id, response_ms, status_code)
        VALUES (%s, %s, %s)
    """, (endpoint_id, response_ms, status_code))


def update_latency_percentiles(cur, endpoint_id, window_hours=24):
    """
    Calculate P50, P95, P99 from the rolling window.
    Uses PostgreSQL's percentile_cont for exact percentiles.
    """
    cur.execute("""
        SELECT 
            COUNT(*) as sample_count,
            percentile_cont(0.50) WITHIN GROUP 
                (ORDER BY response_ms) AS p50,
            percentile_cont(0.95) WITHIN GROUP 
                (ORDER BY response_ms) AS p95,
            percentile_cont(0.99) WITHIN GROUP 
                (ORDER BY response_ms) AS p99
        FROM endpoint_latency
        WHERE endpoint_id = %s
          AND measured_at >= NOW() - INTERVAL '%s hours'
    """, (endpoint_id, window_hours))
    
    row = cur.fetchone()
    
    if row and row['sample_count'] > 0:
        cur.execute("""
            UPDATE endpoint_health
            SET p50_ms = %s,
                p95_ms = %s,
                p99_ms = %s,
                sample_count = %s
            WHERE endpoint_id = %s
        """, (
            int(row['p50']), int(row['p95']), 
            int(row['p99']), row['sample_count'],
            endpoint_id
        ))
    
    return row


def get_slow_endpoints(cur, p95_threshold_ms=2000):
    """
    Find endpoints whose P95 latency exceeds the threshold.
    These are candidates for investigation or circuit opening.
    """
    cur.execute("""
        SELECT endpoint_id, endpoint_url, 
               p50_ms, p95_ms, p99_ms, sample_count,
               state, consecutive_failures
        FROM endpoint_health
        WHERE p95_ms > %s
          AND sample_count >= 20  -- need sufficient samples
        ORDER BY p95_ms DESC
    """, (p95_threshold_ms,))
    
    return cur.fetchall()

PostgreSQL percentile_cont Это агрегатная функция для упорядоченных множеств, которая вычисляет точные процентили. Для больших наборов данных следует перейти к percentile_disc (который возвращает фактическое наблюдаемое значение, а не интерполирует) или используйте приближение t-дайджеста. Для доставки веб-хуков с 24-часовым окном точное вычисление процентилей по исходным данным достаточно быстро до обработки примерно 100 000 измерений на конечную точку.

get_slow_endpoints Эта функция запускается по расписанию каждые 15 минут. Конечные точки с P95 более 2 секунд помечаются для дальнейшего расследования. Для конечных точек с P95 более 5 секунд порог срабатывания автоматического выключателя снижается — им разрешается меньше последовательных сбоев до размыкания цепи, поскольку каждая неудачная доставка занимает рабочий поток на весь период ожидания.

Мониторинг работоспособности доставки веб-хуков

Вот запрос мониторинга, который я запускаю каждые пять минут. Он выдает сводку состояния всего конвейера доставки в одной строке:

SQL

SELECT
    -- Queue depth
    COUNT(*) FILTER (WHERE status = 'pending') AS pending,
    COUNT(*) FILTER (WHERE status = 'failed') AS awaiting_retry,
    COUNT(*) FILTER (WHERE status = 'in_flight') AS in_flight,
    COUNT(*) FILTER (WHERE status = 'dead_letter') AS dead_letter,
    
    -- Delivery rate (last hour)
    COUNT(*) FILTER (
        WHERE status = 'delivered' 
        AND delivered_at >= NOW() - INTERVAL '1 hour'
    ) AS delivered_last_hour,
    
    -- Failure rate (last hour)
    COUNT(*) FILTER (
        WHERE status IN ('failed', 'dead_letter')
        AND created_at >= NOW() - INTERVAL '1 hour'
    ) AS failed_last_hour,
    
    -- Oldest undelivered
    MIN(created_at) FILTER (
        WHERE status IN ('pending', 'failed')
    ) AS oldest_pending,
    
    -- Average delivery latency (last hour, successful only)
    AVG(last_response_ms) FILTER (
        WHERE status = 'delivered'
        AND delivered_at >= NOW() - INTERVAL '1 hour'
    ) AS avg_delivery_ms_last_hour

FROM webhook_queue;

oldest_pending Значение параметра является наиболее важной метрикой в ​​этом запросе. Если оно превышает максимальное окно повторных попыток (сумму всех задержек), значит, что-то не так со структурой — либо рабочий процесс завис, либо конечная точка заблокирована, либо очередь растет быстрее, чем вы можете ее опустошить.

Я получаю оповещения при трех условиях: увеличение количества недоставленных писем (конечные точки постоянно выходят из строя, и никто не занимается расследованием), превышение 30 минут в очереди ожидания (отставание в доставке) и Задержка P95 на каждую конечную точку превышает пороговые значения, указывающие на ухудшение надежности доставки. Третий фактор — это сигнал раннего предупреждения: задержка увеличивается до того, как начнутся сбои. Конечная точка, которая отвечала за 200 мс, а теперь начинает отвечать за 3 секунды, вот-вот начнет выдавать ошибку тайм-аута.

Очередь невостребованных писем — это не просто хранилище.

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

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

питон

def analyze_dead_letters(cur, hours=24):
    """
    Analyze recent dead letter entries for patterns.
    Returns per-endpoint failure analysis.
    """
    cur.execute("""
        SELECT 
            endpoint_id,
            endpoint_url,
            COUNT(*) AS dead_count,
            
            -- Most common error
            MODE() WITHIN GROUP (ORDER BY last_error) AS primary_error,
            
            -- Most common status code
            MODE() WITHIN GROUP (ORDER BY last_status_code) 
                AS primary_status_code,
            
            -- Timing
            MIN(created_at) AS first_dead,
            MAX(created_at) AS last_dead,
            
            -- Average attempts before giving up
            AVG(attempt_count)::INTEGER AS avg_attempts
            
        FROM webhook_queue
        WHERE status = 'dead_letter'
          AND created_at >= NOW() - INTERVAL '%s hours'
        GROUP BY endpoint_id, endpoint_url
        ORDER BY dead_count DESC
        LIMIT 20
    """, (hours,))
    
    return cur.fetchall()

При проверке невостребованных писем я ищу три закономерности.

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

Шаблоны кодов состояния: Резкое увеличение количества недоставленных писем 401/403 означает, что конечная точка сменила учетные данные, и никто не обновил конфигурацию веб-перехватчика. Резкое увеличение количества писем 429 (Слишком много запросов) означает, что вы превышаете лимит запросов и необходимо ограничить скорость соединения.

Постепенное накопление: 2-3 неработающих письма в день для одной и той же конечной точки, равномерно распределенные. Это самая коварная схема — конечная точка в основном работает, но периодически возникают сбои, которые со временем исчерпывают все возможности повторной попытки. Решение обычно заключается в увеличении числа проблем. max_attempts для конкретной конечной точки или для сокращения времени ожидания.

Управление рабочим

Главный цикл, который связывает всё воедино:

питон

import signal
import sys

running = True

def shutdown_handler(signum, frame):
    global running
    running = False
    print(f"Received signal {signum}, shutting down gracefully...")

signal.signal(signal.SIGTERM, shutdown_handler)
signal.signal(signal.SIGINT, shutdown_handler)

def main():
    print("Webhook delivery worker starting...")
    
    while running:
        try:
            deliver_webhooks(batch_size=50)
        except Exception as e:
            print(f"Worker error: {e}")
            time.sleep(5)  # back off on errors
            continue
        
        # Update latency stats every 100 iterations
        # (cheap operation, doesn't need to run every loop)
        if int(time.time()) % 100 == 0:
            conn = get_connection()
            cur = conn.cursor(cursor_factory=RealDictCursor)
            try:
                cur.execute(
                    "SELECT DISTINCT endpoint_id FROM endpoint_health"
                )
                for row in cur.fetchall():
                    update_latency_percentiles(cur, row['endpoint_id'])
                conn.commit()
            finally:
                cur.close()
                conn.close()
        
        # Poll interval — 500ms keeps latency low without
        # hammering the database
        time.sleep(0.5)

    print("Worker shut down cleanly.")

if __name__ == '__main__':
    main()

SIGTERM Обработчик необходим для корректного завершения работы в контейнеризированных средах. Когда Kubernetes отправляет сигнал SIGTERM, рабочий процесс завершает текущую партию, фиксирует транзакцию и завершает работу. Без него строки могут застревать в памяти. in_flight статус, который не обрабатывается никаким работником.

Чего эта система не делает (и когда вам это нужно)

Данная реализация обрабатывает до 10 000 доставок в минуту на одном экземпляре PostgreSQL с 2-3 рабочими процессами. Для большего объема потребуется внести три изменения.

Во-первых, замените очередь PostgreSQL на Redis Streams или RabbitMQ. SELECT FOR UPDATE SKIP LOCKED Такой шаблон создает конкуренцию за запись в таблицу очередей при высокой пропускной способности. Выделенный брокер сообщений устраняет эту проблему.

Во-вторых, добавьте ограничение скорости для каждой конечной точки. Некоторые принимающие конечные точки имеют ограничения скорости (100 запросов в минуту, 1,000 в час). Без ограничения скорости на стороне клиента вы превысите их квоту и получите ошибку 429. Внедрите отдельный токен-бак для каждой конечной точки.

В-третьих, добавьте подписание запросов. Подписи HMAC-SHA256 на полезной нагрузке позволяют принимающей конечной точке проверить, что веб-перехватчик поступил из вашей системы и не был изменен во время передачи. Это обязательное условие для любой системы веб-перехватчиков, отправляющей финансовые данные.

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

Наиболее важный аспект, который большинство команд упускают из виду: измерение самой системы доставки. Если вы не можете ответить на вопрос «какова задержка доставки P95 до конечной точки X за последние 24 часа», вы работаете вслепую. Сначала создайте инструментарий. Все остальное придет потом.

Часто задаваемые вопросы о доставке веб-хуков

Должен ли отправитель веб-хука гарантировать доставку ровно один раз?

Обычно нет. Отправитель должен сделать повторные попытки видимыми и предоставить стабильный идентификатор события; получатель должен обеспечить идемпотентность обработки, чтобы событие могло быть доставлено более одного раза без дублирования бизнес-результата.

Какие неудачные попытки следует повторить?

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

Когда команде следует отказаться от использования очереди, основанной на базе данных?

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

Предыдущая

Лучшее программное обеспечение для отслеживания партнерских программ в сфере iGaming в 2026 году

На следующую

Искусственный интеллект в онлайн-играх: примеры использования ChatGPT в казино.

Цезарь Фиксон
Автор:

Цезарь Фиксон

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

Запросить демо
ШАГ 1 3 г.
Спасибо — вы в очереди.
Специалист компании NowG свяжется с вами в течение одного рабочего дня, чтобы согласовать время осмотра помещения.
Индекс