Bygge et Webhook-leveringssystem: Pålitelige tilbakesendinger i 2026

En praktisk guide for 2026 til holdbar webhook- og affiliate-postback-levering med køer, jitterede forsøk, sikringsbrytere, idempotens, signaturer og telemetri for latenstid per endepunkt.
webhook-leveringssystem - Bygge et webhook-leveringssystem: Pålitelige tilbakesendinger i 2026

Sist oppdatert 24. juni 2026 av Cæsar Fikson

Direkte svar: En pålitelig webhook- eller affiliate-postback-leveringstjeneste trenger en holdbar kø, begrensede nye forsøk med jitter, kretsbryting per endepunkt, duplikatbeskyttelse, signaturverifisering og telemetri som viser om leveringen går tregere før den mislykkes.

Denne implementeringen holder utførelsesmodellen bevisst liten: Python-arbeidere og PostgreSQL for en kø, leveringshistorikk og endepunktshelse. Det er et praktisk utgangspunkt for B2B SaaS, iGaming-tilknyttede operasjoner og ethvert produkt som må levere utgående hendelser uten å behandle en enkelt HTTP-timeout som en tapt konvertering.

Hva leveringskontrakten må garantere

Kontroll: Hvorfor det betyr noe Driftssjekk
Holdbar hendelseshistorikk Hendelser overlever omstart av arbeidere. Hver akseptert levering har en stabil ID og status.
Signert forespørsel Mottakere kan bekrefte avsenderen og oppdage kroppsmanipulering. Bruk en tidsstemplet signatur og avvis foreldede forespørsler.
Duplikatbeskyttelse Ellers kan nye forsøk føre til dupliserte konverteringer eller oppdateringer. Send en hendelses-ID og gjør mottakeren idempotent.
Avgrensede nye forsøk Midlertidige feil gjenopprettes uten å overbelaste et degradert endepunkt. Trekk av ved jitter; stopp etter en dokumentert forsøksgrense.
Endepunkttelemetri Bare kødybden skjuler trege eller sviktende partnere. Spor suksessrate, eldste ventende hendelse og latenspersentiler per endepunkt.

Sikkerhets- og duplikatkontroller kommer før ny justering

Behandle den utgående forespørselsteksten som sensitive driftsdata. Signer den nøyaktige forespørselsteksten, inkluder et leveringstidsstempel og en uforanderlig hendelses-ID, roter signeringshemmeligheter og sørg for at mottakerapplikasjonen trygt ignorerer en avspilling av den samme hendelsen. Et vellykket HTTP-svar er ikke bevis på at en forretningshendelse ble brukt nøyaktig én gang; mottakersiden må ta den avgjørelsen.

For en arbeidsflyt for postback for tilknyttede selskaper eller iGaming, hold klikk-ID-er, konverterings-ID-er og utbetalingsrelevante felt unna loggene med mindre tilgangskontroller og oppbevaringsregler er eksplisitte. Den eksisterende Scaleo-referansen nedenfor er kun relevant som et eksempel på latens for postback for tilknyttede selskaper; den erstatter ikke å dokumentere leveringskontrakten mellom dine egne tjenester.

Nyttige implementeringsreferanser: Verifisering av Stripe webhook-signatur, PostgreSQL SELECT og SKIP LÅSTog AWS-veiledning om eksponentiell tilbakeslag og jitter.

Datamodellen

Alt starter med to tabeller: én for leveringskøen og én for latensmålingene.

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

Tre ting å merke seg om dette skjemaet.

Først webhook_queue tabellen bruker en next_attempt_at kolonne i stedet for en separat planleggingsmekanisme. Arbeideren avspør etter rader der status IN ('pending', 'failed') AND next_attempt_at <= NOW()Dette er en fattigmannskø, og den fungerer fint opptil omtrent 10 000 leveringer per minutt. Utover det, bytt til en skikkelig meldingsmegler.

Det andre, endpoint_latency Tabellen fungerer som en ringbuffer. Jeg sletter jevnlig rader som er eldre enn 24 timer. Latensperioden i endpoint_health beregnes fra dette rullerende vinduet – de representerer nylig atferd, ikke historiske gjennomsnitt.

Tredje, endpoint_health Tabellen implementerer tilstandsmaskinen for effektbryteren. Mer om dette nedenfor.

Leveringsarbeideren

Kjernearbeiderløkken er bevisst enkel. Kompleksitet hører hjemme i gjentakelseslogikken og sikringsbryteren, ikke i selve leveringsbanen.

python

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

Ocuco SELECT FOR UPDATE SKIP LOCKED klausulen er kritisk for å kjøre flere arbeiderinstanser. Uten SKIP LOCKED, ville to arbeidere blokkere på samme rad. Med den henter hver arbeider en annen gruppe med ventende webhooks. Dette gir deg horisontal skalering ved ganske enkelt å starte flere arbeiderprosesser.

Ocuco time.monotonic() ring i stedet for time.time() er bevisst. time.time() kan hoppe bakover under NTP-justeringer. time.monotonic() går aldri bakover, noe som er viktig når du måler latens på under et sekund.

Prøv logikk på nytt med eksponensiell backoff og jitter

Når en levering mislykkes, avgjør timingen av nytt forsøk om systemet gjenoppretter seg uten problemer eller om det skaper en tordnende flokk som hamrer ned et endepunkt som sliter.

python

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)

Hvorfor full jitter i stedet for dekorrelatert jitter eller lik jitter? AWS publiserte den definitive analysen av dette. Full jitter (randomisering mellom 0 og den eksponensielle grensen) gir den laveste totale fullføringstiden på tvers av alle klienter. Lik jitter (randomisering mellom halve grensen og full grense) er mer konservativt, men tregere for å tømme etterslepet av nye forsøk. For webhook-levering der du har mange uavhengige endepunkter, er full jitter det riktige valget fordi hvert endepunkts nye forsøk er uavhengige – du koordinerer ikke mellom dem.

Sikkerhetsbryter: Slutt å hamre på ødelagte endepunkter

Brytermønsteret hindrer systemet i å kaste bort ressurser på endepunkter som stadig svikter. Uten det akkumulerer et dødt endepunkt hundrevis av ventende nye forsøk som alle får tidsavbrudd på 15 sekunder hver – og bruker dermed opp arbeidskapasiteten din på leveranser som aldri vil lykkes.

python

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

Halvåpen → åpen eskalering med dobbel nedkjølingstid er detaljen de fleste implementeringer går glipp av. Hvis et endepunkt feiler under proben (halvåpen tilstand), bør du ikke prøve på nytt om 5 minutter. Endepunktet er fortsatt ødelagt. Doble nedkjølingstiden til 10 minutter, deretter 20, med en grense på 1 time. Dette forhindrer at effektbryteren blir en periodisk hammermekanisme.

Latenspersentilsporing

Gjennomsnittene lyver. Et endepunkt med en gjennomsnittlig responstid på 200 ms kan reagere på 50 ms i 95 % av tilfellene og 3,000 ms de resterende 5 %. Gjennomsnittet ser fint ut. P95 avslører et problem som påvirker 1 av 20 leveranser.

python

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-er percentile_cont er en ordnet mengde aggregatfunksjon som beregner eksakte persentiler. For store datasett bytter du til percentile_disc (som returnerer en faktisk observert verdi i stedet for interpolering) eller bruk en t-digest-tilnærming. For webhook-levering med et 24-timers vindu er eksakte persentiler på rådataene raske nok opptil omtrent 100 000 målinger per endepunkt.

Ocuco get_slow_endpoints -funksjonen er det jeg kjører som en planlagt sjekk hvert 15. minutt. Endepunkter med P95 over 2 sekunder blir flagget for undersøkelse. Endepunkter med P95 over 5 sekunder får redusert terskelverdien for kretsbryter – de får færre påfølgende feil før kretsen åpnes, fordi hver mislykkede levering binder opp en arbeidstråd for hele tidsavbruddsvarigheten.

Overvåking av leveringstilstanden til webhooken

Her er overvåkingsspørringen jeg kjører hvert femte minutt. Den produserer et helsesammendrag på én rad av hele leveringsprosessen:

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;

Ocuco oldest_pending Verdien er den viktigste metrikken i denne spørringen. Hvis den er eldre enn det maksimale vinduet for nye forsøk (summen av alle tilbakekallingsforsinkelser), er det noe strukturelt galt – enten sitter arbeideren fast, endepunktet er i et blackhole, eller køen vokser raskere enn du kan tømme den.

Jeg varsler under tre forhold: antall døde brev øker (endepunktene svikter permanent og ingen undersøker), ventende kø som overstiger 30 minutter (leveringen henger etter), og P95-forsinkelse per endepunkt som overstiger terskler som indikerer redusert leveringspålitelighet Det tredje er det tidlige varslingssignalet – latensen øker før feil gjør det. Et endepunkt som svarte på 200 ms og begynner å svare på 3 sekunder, er i ferd med å begynne å få tidsavbrudd.

Dødbrevkøen er ikke bare lagring

De fleste team implementerer en «dead letter»-kø som en tabell der mislykkede webhooks havner for å dø. De sjekker den av og til under hendelsesrespons. Dette er bortkastet arbeid.

Køen med døde bokstaver er ditt mest verdifulle feilsøkingsdatasett. Hver rad representerer en levering som systemet ditt prøvde flere ganger og ga opp. Mønsteret med døde bokstaver forteller deg ting som suksessmålingene aldri vil gjøre.

python

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

Når jeg gjennomgår døde bokstaver, ser jeg etter tre mønstre.

Klyngefeil: 50 døde bokstaver for samme endepunkt i samme time betyr at endepunktet gikk ned og ikke gjenopprettet seg innen forsøksvinduet. Handling: forleng forsøksvinduet eller implementer manuell ny kø.

Statuskodemønstre: En topp i 401/403 døde bokstaver betyr at endepunktet roterte legitimasjonsinformasjonen og at ingen oppdaterte webhook-konfigurasjonen. En topp i 429 (For mange forespørsler) betyr at du overskrider hastighetsgrensen deres og må begrense den.

Gradvis akkumulering: 2–3 døde bokstaver per dag for et enkelt endepunkt, jevnt fordelt. Dette er det mest utspekulerte mønsteret – endepunktet fungerer stort sett, men har periodiske feil som utmatter nye forsøk over tid. Løsningen øker vanligvis. max_attempts for det spesifikke endepunktet eller å redusere tidsavbruddet.

Kjører arbeideren

Hovedsløyfen som binder alt sammen:

python

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

Ocuco SIGTERM Behandleren er viktig for rene nedstengninger i containeriserte miljøer. Når Kubernetes sender SIGTERM, fullfører arbeideren sin gjeldende batch, utfører transaksjonen og avslutter. Uten dette får du rader som sitter fast i in_flight status uten at noen arbeider behandler dem.

Hva dette systemet ikke gjør (og når du trenger mer)

Denne implementeringen håndterer opptil omtrent 10 000 leveranser per minutt på en enkelt PostgreSQL-instans med 2–3 arbeidsprosesser. Utover det er det behov for tre endringer.

Først, erstatt PostgreSQL-køen med Redis Streams eller RabbitMQ. SELECT FOR UPDATE SKIP LOCKED Mønsteret skaper skrivekonflikt på køtabellen ved høy gjennomstrømning. En dedikert meldingsmegler eliminerer dette.

For det andre, legg til hastighetsbegrensning per endepunkt. Noen mottakende endepunkter har hastighetsgrenser (100 forespørsler per minutt, 1,000 per time). Uten hastighetsbegrensning på klientsiden vil du bruke opp kvoten deres og få 429. Implementer en token-bøtte per endepunkt.

For det tredje, legg til forespørselssignering. HMAC-SHA256-signaturer på nyttelasten lar det mottakende endepunktet bekrefte at webhooken kom fra systemet ditt og ikke ble tuklet med under overføring. Dette er tabellinnsatser for ethvert webhook-system som sender økonomiske data.

Systemet i denne artikkelen er fundamentet. Det håndterer de vanskelige problemene – logikk for nye forsøk, kretsbryting, latensmåling, analyse av døde bokstaver – som alle webhook-leveringssystemer trenger uavhengig av skala. De spesifikke komponentene du legger til (meldingsmegler, hastighetsbegrenser, forespørselssignering) avhenger av gjennomstrømnings- og sikkerhetskravene dine.

Den delen som betyr mest er den delen de fleste team hopper over: å måle selve leveringssystemet. Hvis du ikke kan svare på «hva er P95-leveringslatensen til endepunkt X de siste 24 timene», opererer du i blinde. Bygg instrumentasjonen først. Alt annet følger.

Vanlige spørsmål om webhook-levering

Bør en webhook-avsender love levering nøyaktig én gang?

Vanligvis ikke. En avsender bør gjøre nye forsøk synlige og oppgi en stabil hendelses-ID; mottakeren bør gjøre behandlingen idempotent slik at en hendelse kan leveres mer enn én gang uten å duplisere forretningsresultatet.

Hvilke feil bør prøves på nytt?

Prøv bare på nytt feil som kontrakten din klassifiserer som forbigående, for eksempel nettverksfeil, tidsavbrudd og valgte serversvar. Ikke prøv gjentatte ganger på nytt feilformede forespørsler, autentiseringsfeil eller andre permanente feil uten en eksplisitt utbedringsvei.

Når bør et team gå videre fra en databasebasert kø?

Flytt når målt konkurranse, alder på ordrebeholdning, gjennomstrømning eller driftsmessige gjenopprettingsbehov viser at databasekøen ikke lenger oppfyller leveringskontrakten. Kapasitetsbeslutninger bør følge observert arbeidsmengde, ikke en generisk påstand om forespørselsrate.

Forrige Artikkel

Beste programvare for sporing av iGaming-affiliate i 2026

Neste Artikkel

AI i iGaming: ChatGPT-brukstilfeller for kasinoer

Cæsar Fikson
Forfatter:

Cæsar Fikson

Jeg er en iGaming-dataanalytiker som spesialiserer seg på å undersøke og tolke data relatert til online spillplattformer og gamblingaktiviteter, samt markedstrender. Jeg analyserer spilleratferd, spillytelse og inntektstrender for å optimalisere spillopplevelser og forretningsstrategier.

Be om en demonstrasjon
TRINN 1 AV 3
Takk – du er i køen.
En NowG-løsningsingeniør vil ta kontakt med deg innen én virkedag for å avtale en gjennomgang.
Index