Vytvoření systému pro doručování webhooků: Spolehlivé postbacky v roce 2026

Praktický průvodce pro rok 2026 pro odolné doručování webhooků a affiliate postbacků s frontami, jitterovanými opakovanými pokusy, jističi, idempotencí, podpisy a telemetrií latence pro jednotlivé koncové body.
systém doručování webhooků - Vytvoření systému doručování webhooků: Spolehlivé postbacky v roce 2026

Naposledy aktualizováno 24. června 2026 do Caesar Fikson

Přímá odpověď: Spolehlivá doručovací služba webhook nebo affiliate-postback potřebuje odolnou frontu, omezené opakované pokusy s jitterem, přerušení okruhu pro každý koncový bod, ochranu proti duplicitě, ověřování podpisu a telemetrii, která ukazuje, zda se doručování zpomaluje, než dojde k jeho selhání.

Tato implementace záměrně udržuje model provádění malý: Python workers a PostgreSQL pro frontu, historii doručování a stav koncových bodů. Je to praktický výchozí bod pro B2B SaaS, affiliate operace iGaming a jakýkoli produkt, který musí doručovat odchozí události bez toho, aby byl jediný timeout HTTP považován za ztracenou konverzi.

Co musí zaručovat dodací smlouva

ovládání Proč je to důležité Provozní kontrola
Trvalý záznam událostí Události přežijí restartování pracovníků. Každá přijatá zásilka má stabilní ID a status.
Podepsaná žádost Příjemci mohou ověřit odesílatele a detekovat neoprávněnou manipulaci s tělem. Používejte podpis s časovým razítkem a odmítejte zastaralé požadavky.
Ochrana proti duplicitě Opakované pokusy mohou jinak vést k duplicitním konverzím nebo aktualizacím. Odešle ID události a nastaví příjemce jako idempotentní.
Omezené opakované pokusy Dočasné selhání se zotaví bez zahlcení degradovaného koncového bodu. Ukončete s jitterem; zastavte po uplynutí zdokumentovaného limitu pokusů.
Telemetrie koncových bodů Už jen hloubka fronty skrývá pomalé nebo neúspěšné partnery. Sledujte míru úspěšnosti, nejstarší čekající událost a percentily latence pro každý koncový bod.

Zabezpečení a duplicitní ovládací prvky mají přednost před laděním opakovaných pokusů

S odchozím tělem zacházejte jako s citlivými provozními daty. Podepište přesné tělo požadavku, uveďte časové razítko doručení a neměnné ID události, rotujte tajné klíče podpisu a zajistěte, aby přijímající aplikace bezpečně ignorovala přehrání stejné události. Úspěšná odpověď HTTP není důkazem, že obchodní událost byla použita přesně jednou; toto rozhodnutí musí učinit přijímající strana.

V případě affiliate nebo iGaming postback workflow uchovávejte ID kliknutí, ID konverzí a pole relevantní pro výplaty mimo protokoly, pokud nejsou explicitně uvedeny kontroly přístupu a pravidla uchovávání. Stávající reference Scaleo níže je relevantní pouze jako příklad latence affiliate postbacku; nenahrazuje dokumentaci smlouvy o dodání mezi vašimi vlastními službami.

Užitečné reference k implementaci: Ověření podpisu webhooku Stripe, PostgreSQL SELECT a SKIP LOCKED, a Pokyny AWS k exponenciálnímu poklesu a chvění.

Datový model

Všechno začíná dvěma tabulkami: jednou pro frontu doručení a druhou pro měření latence.

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
);

Tři věci, které je třeba si u tohoto schématu uvědomit.

Za prvé, webhook_queue tabulka používá next_attempt_at sloupec namísto samostatného mechanismu plánování. Pracovník se dotazuje na řádky, kde status IN ('pending', 'failed') AND next_attempt_at <= NOW()Tohle je fronta zpoždění pro chudé a funguje dobře až do zhruba 10 000 doručení za minutu. Kromě toho použijte řádného zprostředkovatele zpráv.

Za druhé, endpoint_latency Tabulka funguje jako kruhová vyrovnávací paměť. Pravidelně mažu řádky starší než 24 hodin. Percentuální hodnoty latence v endpoint_health se počítají z tohoto posuvného okna – představují nedávné chování, nikoli historické průměry.

Za třetí, endpoint_health Tabulka implementuje stavový automat jističe. Více o tom níže.

Doručovatel

Základní pracovní smyčka je záměrně jednoduchá. Složitost patří do logiky opakování a jističe, nikoli do samotné cesty doručení.

krajta

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]}'
        }

Jedno SELECT FOR UPDATE SKIP LOCKED Klauzule je kritická pro spouštění více instancí pracovních procesů. Bez SKIP LOCKED, dva workeři by blokovali stejný řádek. Díky tomu každý worker nabere jinou dávku čekajících webhooků. To vám dává horizontální škálování pouhým spuštěním dalších worker procesů.

Jedno time.monotonic() volat místo time.time() je úmyslné. time.time() může během úprav NTP skočit zpět. time.monotonic() nikdy se nevrací zpět, což je důležité při měření latence v délce kratší než jedna sekunda.

Logika opakování s exponenciálním odkladem a jitterem

Když se doručení nezdaří, načasování opakování určuje, zda se systém elegantně obnoví, nebo vytvoří hromové stádo, které buší do problémového koncového bodu.

krajta

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)

Proč zvolit plný jitter místo dekorelovaného jitteru nebo ekvivalentního jitteru? Společnost AWS publikovala definitivní analýzu tohoto tématu. Plný jitter (náhodné přepínání mezi 0 a exponenciálním limitem) vede k nejnižší celkové době dokončení u všech klientů. Equal jitter (náhodné přepínání mezi polovinou limitu a plným limitem) je konzervativnější, ale pomalejší pro odčerpávání backlogu opakovaných pokusů. Pro doručování webhooků, kde máte mnoho nezávislých koncových bodů, je plný jitter správnou volbou, protože opakované pokusy každého koncového bodu jsou nezávislé – nekoordinujete je mezi sebou.

Jistič: Přestaňte bušit do poškozených koncových bodů

Vzor jističe zabraňuje plýtvání zdrojů systémem na koncových bodech, které neustále selhávají. Bez něj nefunkční koncový bod hromadí stovky čekajících pokusů, které všechny po 15 sekundách vyprší – a tím spotřebovává vaši pracovní kapacitu na doručení, která nikdy nebudou úspěšná.

krajta

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,))

Eskalace half_open → open s dvojnásobnou dobou ochlazování je detail, který většina implementací opomíjí. Pokud koncový bod selže během testování (stav half_open), nechcete to zkoušet znovu za dalších 5 minut. Koncový bod je stále nefunkční. Zdvojnásobte dobu ochlazování na 10 minut, poté na 20 minut, s omezením na 1 hodinu. Tím se zabrání tomu, aby se jistič stal periodickým úderovým mechanismem.

Sledování percentilu latence

Průměry lžou. Koncový bod s průměrnou dobou odezvy 200 ms může reagovat za 50 ms v 95 % případů a za 3 000 ms ve zbývajících 5 %. Průměr vypadá dobře. P95 odhaluje problém, který postihuje 1 z 20 doručení.

krajta

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 je agregační funkce s uspořádanými množinami, která počítá přesné percentily. Pro velké datové sady byste přepnuli na percentile_disc (který vrací skutečnou pozorovanou hodnotu, nikoli interpolaci) nebo použijte aproximaci t-digest. Pro doručování webhooků s 24hodinovým oknem jsou přesné percentily na nezpracovaných datech dostatečně rychlé až do přibližně 100 000 měření na koncový bod.

Jedno get_slow_endpoints Funkce je to, co spouštím jako plánovanou kontrolu každých 15 minut. Koncové body s P95 nad 2 sekundy jsou označeny k prošetření. Koncovým bodům s P95 nad 5 sekund se snižuje prahová hodnota jističe – je jim povoleno méně po sobě jdoucích selhání, než se okruh otevře, protože každé neúspěšné doručení váže pracovní vlákno po celou dobu časového limitu.

Monitorování stavu doručování webhooků

Zde je monitorovací dotaz, který spouštím každých pět minut. Vygeneruje jednořádkové shrnutí stavu celého doručovacího kanálu:

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;

Jedno oldest_pending Hodnota je nejdůležitější metrikou v tomto dotazu. Pokud je starší než maximální doba opakování (součet všech zpoždění backoffu), je něco strukturálně špatně – buď je worker zaseknutý, koncový bod je blackhole, nebo fronta roste rychleji, než ji stihnete vyprázdnit.

Zobrazuji upozornění na tři stavy: rostoucí počet nedoručených zpráv (koncové body trvale selhávají a nikdo to nezkoumá), stáří čekající fronty přesahující 30 minut (doručení se zpožďuje) a Latence P95 na koncový bod překračující prahové hodnoty, které indikují sníženou spolehlivost doručování Třetím je signál včasného varování – latence se zvyšuje ještě před selháním. Koncový bod, který reagoval 200 ms a začne reagovat za 3 sekundy, se brzy začne prodlužovat.

Fronta mrtvého dopisu není jen úložiště

Většina týmů implementuje frontu nefunkčních odkazů jako tabulku, kam se neúspěšné webhooky ukládají. Občas ji kontrolují během reakce na incidenty. To je plýtvání.

Fronta mrtvých dopisů je vaše nejcennější ladicí datová sada. Každý řádek představuje doručení, o které se váš systém několikrát pokusil a nakonec se vzdal. Vzor mrtvých dopisů vám říká věci, které vám metriky úspěšnosti nikdy neřeknou.

krajta

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()

Když procházím mrtvé dopisy, hledám tři vzorce.

Selhání clusteru: 50 mrtvých zpráv pro stejný koncový bod ve stejné hodině znamená, že koncový bod selhal a neobnovil se v rámci okna pro opakování. Akce: prodlužte okno pro opakování nebo implementujte ruční opětovné zařazení do fronty.

Vzory stavových kódů: Nárůst chyb 401/403 (mrtvé zprávy) znamená, že koncový bod změnil přihlašovací údaje a nikdo neaktualizoval konfiguraci webhooku. Nárůst chyb 429 (příliš mnoho požadavků) znamená, že překračujete jejich limit rychlosti a je třeba omezit rychlost.

Postupná akumulace: 2–3 nefunkční zprávy denně pro jeden koncový bod, rovnoměrně rozložené. Toto je nejzáludnější vzorec – koncový bod většinou funguje, ale občas se vyskytují chyby, které časem vyčerpávají opakované pokusy. Oprava se obvykle zvyšuje. max_attempts pro daný koncový bod nebo zkrácení časového limitu.

Spuštění pracovníka

Hlavní smyčka, která vše spojuje:

krajta

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()

Jedno SIGTERM Obslužná rutina je nezbytná pro čisté ukončení v kontejnerizovaných prostředích. Když Kubernetes odešle SIGTERM, worker dokončí svou aktuální dávku, potvrdí transakci a ukončí práci. Bez ní se řádky zasekávají v in_flight stav bez pracovníka, který je zpracovává.

Co tento systém neumí (a kdy potřebujete víc)

Tato implementace zvládá až přibližně 10 000 doručení za minutu na jedné instanci PostgreSQL s 2–3 pracovními procesy. Kromě toho jsou potřeba tři změny.

Nejprve nahraďte frontu PostgreSQL streamy Redis nebo RabbitMQ. SELECT FOR UPDATE SKIP LOCKED Tento vzorec vytváří konflikt zápisu v tabulce fronty při vysoké propustnosti. Vyhrazený zprostředkovatel zpráv to eliminuje.

Za druhé, přidejte omezení rychlosti pro každý koncový bod. Některé přijímající koncové body mají omezení rychlosti (100 požadavků za minutu, 1 000 za hodinu). Bez omezení rychlosti na straně klienta překročíte jejich kvótu a dostanete chybu 429. Implementujte token bucket pro každý koncový bod.

Za třetí, přidejte podepisování požadavků. Podpisy HMAC-SHA256 v datové části umožňují přijímajícímu koncovému bodu ověřit, zda webhook pochází z vašeho systému a nebyl během přenosu pozměněn. Toto je klíčové pro jakýkoli systém webhooků, který odesílá finanční data.

Systém popsaný v tomto článku je základem. Řeší složité problémy – logiku opakování, přerušení okruhu, měření latence, analýzu nedoručených zpráv – které každý systém doručování webhooků potřebuje bez ohledu na rozsah. Konkrétní komponenty, které nainstalujete (zprostředkovatel zpráv, omezovač rychlosti, podepisování požadavků), závisí na vašich požadavcích na propustnost a zabezpečení.

Část, na které záleží nejvíc, je ta, kterou většina týmů přeskakuje: měření samotného systému doručování. Pokud nedokážete odpovědět na otázku „jaká je latence doručení P95 do koncového bodu X za posledních 24 hodin“, pracujete naslepo. Nejprve sestavte instrumentaci. Všechno ostatní pak následuje.

Nejčastější dotazy k doručování webhooků

Měl by odesílatel webhooku slibovat doručení přesně jednou?

Obvykle ne. Odesílatel by měl zviditelnit opakované pokusy a poskytnout stabilní ID události; příjemce by měl nastavit procesní idempotent, aby událost mohla být doručena vícekrát, aniž by došlo k duplikaci obchodního výsledku.

Které selhání by se mělo zopakovat?

Opakujte pouze selhání, která vaše smlouva klasifikuje jako přechodná, například chyby sítě, časové limity a vybrané odpovědi serveru. Neopakujte opakovaně chybné požadavky, selhání ověřování nebo jiné trvalé chyby bez explicitní cesty k nápravě.

Kdy by se měl tým přesunout za hranice fronty založené na databázi?

Přesunout se, když naměřené konflikty, stáří nevyřízených objednávek, propustnost nebo potřeby obnovy provozu ukážou, že databázová fronta již nesplňuje požadavky kontraktu na dodání. Rozhodnutí o kapacitě by se měla řídit pozorovanou pracovní zátěží, nikoli obecným tvrzením o míře požadavků.

Předchozí článek

Nejlepší software pro sledování affiliate partnerů v iGamingu v roce 2026

Další článek

AI v iGamingu: Případy použití ChatGPT pro kasina

Caesar Fikson
Autor:

Caesar Fikson

Jsem iGaming datový analytik specializující se na zkoumání a interpretaci dat týkajících se online herních platform a hazardních aktivit, jakož i tržních trendů. Analyzuji chování hráčů, herní výkon a trendy tržeb s cílem optimalizovat herní zážitky a obchodní strategie.

Žádost o demo
KROK 1 Z 3
Díky – jste ve frontě.
Technik řešení NowG se s vámi do jednoho pracovního dne spojí a domluví si s vámi konzultaci.
index