Створення системи доставки вебхуків: надійні постбеки у 2026 році

Практичний посібник 2026 року щодо надійної доставки вебхуків та партнерських постбеків з урахуванням черг, джиттерованих повторних спроб, автоматичних вимикачів, ідемпотентності, підписів та телеметрії затримки для кожної кінцевої точки.
система доставки вебхуків - Створення системи доставки вебхуків: надійні постбеки у 2026 році

Останнє оновлення 24 червня 2026 року Цезар Фіксон

Пряма відповідь: Надійний вебхук або сервіс доставки партнерського постбеку потребує стійкої черги, обмежених повторних спроб з джиттером, розриву ланцюга для кожної кінцевої точки, захисту від дублікатів, перевірки підписів і телеметрії, яка показує, чи сповільнюється доставка, перш ніж вона станеться збій.

Ця реалізація навмисно зберігає модель виконання малою: Python workers та PostgreSQL для черги, історії доставки та справності кінцевих точок. Це практична відправна точка для B2B SaaS, партнерських операцій iGaming та будь-якого продукту, який повинен доставляти вихідні події без обробки жодного тайм-ауту HTTP як втраченої конверсії.

Що має гарантувати договір поставки

Контроль Чому це важливо? Операційна перевірка
Довготривалий запис подій Події зберігаються після перезапуску працівників. Кожна прийнята доставка має стабільний ідентифікатор та статус.
Підписаний запит Одержувачі можуть перевірити відправника та виявити втручання в тіло. Використовуйте підпис із позначкою часу та відхиляйте застарілі запити.
Захист від дублікатів В іншому випадку повторні спроби можуть призвести до дублікатів конверсій або оновлень. Надіслати ідентифікатор події та зробити одержувач ідемпотентним.
Обмежені повторні спроби Тимчасові збої відновлюються без перевантаження деградованої кінцевої станції. Зупиніться після задокументованого обмеження спроб, якщо спостерігається тремтіння.
Телеметрія кінцевої точки Сама лише глибина черги приховує повільних або невдалих партнерів. Відстежуйте рівень успішності, найстарішу подію, що очікує, та процентилі затримки для кожної кінцевої точки.

Безпека та дублікати елементів керування передують повторній спробі налаштування

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

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

Корисні посилання на впровадження: Перевірка підпису вебхука Stripe, PostgreSQL SELECT та 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,))

Ескалація half_open → open з подвоєним часом відновлення — це деталь, яку більшість реалізацій пропускають. Якщо кінцева точка дає збій під час зондування (стан half_open), ви не хочете повторювати спробу через 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 році

Наступна стаття

Штучний інтелект в iGaming: варіанти використання ChatGPT для казино

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

Цезар Фіксон

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

Запросити демо
КРОК 1 З 3
Дякую — ви в черзі.
Інженер з рішень NowG зв'яжеться з вами протягом одного робочого дня, щоб запланувати ознайомлення.
індекс