Posledná aktualizácia 24. júna 2026 do Caesar Fikson
Priama odpoveď: Spoľahlivá doručovacia služba webhook alebo affiliate-postback potrebuje odolný front, ohraničené opakovania s jitterom, prerušenie okruhu pre každý koncový bod, ochranu pred duplikátmi, overenie podpisu a telemetriu, ktorá ukazuje, či sa doručovanie spomaľuje skôr, ako zlyhá.
Táto implementácia zámerne udržiava model vykonávania malý: Python workers a PostgreSQL pre front, históriu doručovania a stav koncových bodov. Je to praktický východiskový bod pre B2B SaaS, affiliate operácie iGaming a akýkoľvek produkt, ktorý musí doručovať odchádzajúce udalosti bez toho, aby sa jediný časový limit HTTP považoval za stratenú konverziu.
Čo musí zaručiť zmluva o dodaní
| ovládanie | Prečo je to dôležité | Prevádzková kontrola |
|---|---|---|
| Trvalý záznam o udalostiach | Udalosti prežijú reštartovanie pracovníkov. | Každá prijatá zásielka má stabilné ID a status. |
| Podpísaná žiadosť | Prijímatelia môžu overiť odosielateľa a odhaliť manipuláciu s telom. | Používajte podpis s časovou pečiatkou a odmietajte zastarané požiadavky. |
| Ochrana pred duplikátmi | Opakované pokusy môžu inak viesť k duplicitným konverziám alebo aktualizáciám. | Odošlite ID udalosti a nastavte prijímača ako idempotentného. |
| Ohraničené opakovania | Prechodné zlyhania sa zotavujú bez toho, aby zahltili degradovaný koncový bod. | Ustúpte pri chvení; zastavte po zdokumentovanom limite pokusov. |
| Telemetria koncových bodov | Samotná hĺbka frontu skrýva pomalých alebo zlyhávajúcich partnerov. | Sledujte mieru úspešnosti, najstaršiu čakajúcu udalosť a percentily latencie pre každý koncový bod. |
Zabezpečenie a duplicitné kontroly majú prednosť pred opakovaným ladením
Odchádzajúce telo žiadosti považujte za citlivé prevádzkové údaje. Podpíšte presné telo požiadavky, uveďte časovú pečiatku doručenia a nemenné ID udalosti, striedajte tajné kľúče podpisovania a zabezpečte, aby prijímajúca aplikácia bezpečne ignorovala opakovanie tej istej udalosti. Úspešná odpoveď HTTP nie je dôkazom, že obchodná udalosť bola použitá presne raz; toto rozhodnutie musí urobiť prijímajúca strana.
V prípade affiliate alebo iGaming postback pracovného postupu uchovávajte ID kliknutí, ID konverzií a polia relevantné pre výplaty mimo protokolov, pokiaľ nie sú explicitne uvedené pravidlá riadenia prístupu a uchovávania údajov. Existujúca referencia Scaleo uvedená nižšie je relevantná iba ako príklad latencie affiliate postbacku; nenahrádza dokumentáciu zmluvy o doručovaní medzi vašimi vlastnými službami.
Užitočné referencie implementácie: Overenie podpisu webhooku Stripe, PostgreSQL SELECT a SKIP LOCKEDa Pokyny AWS týkajúce sa exponenciálneho poklesu a jitteru.
Dátový model
Všetko začína dvoma tabuľkami: jednou pre doručovací front a druhou pre merania latencie.
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
);
Tri veci, ktoré si treba všimnúť o tejto schéme.
Po prvé, webhook_queue tabuľka používa next_attempt_at stĺpec namiesto samostatného mechanizmu plánovania. Pracovník sa pýta na riadky, kde status IN ('pending', 'failed') AND next_attempt_at <= NOW()Toto je len front oneskorení pre chudobných a funguje dobre až do približne 10 000 doručení za minútu. Okrem toho použite vhodného sprostredkovateľa správ.
Po druhé, endpoint_latency Tabuľka funguje ako kruhová vyrovnávacia pamäť. Pravidelne čistím riadky staršie ako 24 hodín. Percentily latencie v endpoint_health sa vypočítavajú z tohto posuvného okna – predstavujú nedávne správanie, nie historické priemery.
Po tretie, endpoint_health tabuľka implementuje stavový automat ističa. Viac o tom nižšie.
Doručovateľ
Základná pracovná slučka je zámerne jednoduchá. Zložitosť patrí do logiky opakovania a ističa, nie do samotnej cesty doručenia.
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]}'
}
SELECT FOR UPDATE SKIP LOCKED Klauzula je kritická pre spúšťanie viacerých inštancií pracovníkov. Bez SKIP LOCKED, dvaja workeri by blokovali ten istý riadok. Vďaka tomu každý worker prevezme inú dávku čakajúcich webhookov. To vám dáva horizontálne škálovanie jednoduchým spustením ďalších worker procesov.
time.monotonic() zavolajte namiesto time.time() je úmyselné. time.time() môže počas úprav NTP skočiť dozadu. time.monotonic() nikdy nejde späť, čo je dôležité pri meraní latencie kratšej ako sekunda.
Logika opakovania s exponenciálnym oddialením a jitterom
Keď doručenie zlyhá, načasovanie opakovania určuje, či sa váš systém elegantne obnoví alebo vytvorí hromové stádo, ktoré udrie 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)
Prečo použiť plný jitter namiesto dekorelovaného jitteru alebo rovnakého jitteru? Spoločnosť AWS publikovala definitívnu analýzu tejto témy. Plný jitter (náhodné prepínanie medzi 0 a exponenciálnym stropom) produkuje najnižší celkový čas dokončenia u všetkých klientov. Rovný jitter (náhodné prepínanie medzi polovičným stropom a plným stropom) je konzervatívnejší, ale pomalší pri odčerpávaní nevybavených pokusov o opakovanie. Pre doručovanie webhookov, kde máte veľa nezávislých koncových bodov, je plný jitter tou správnou voľbou, pretože opakované pokusy každého koncového bodu sú nezávislé – nekoordinujete ich medzi sebou.
Istič: Prestaňte búchať do poškodených koncových bodov
Vzor ističa zabraňuje plytvaniu zdrojmi systému na koncových bodoch, ktoré neustále zlyhávajú. Bez neho nefunkčný koncový bod hromadí stovky čakajúcich pokusov, ktoré všetky vypršia po 15 sekundách – čím sa spotrebúva kapacita vašich pracovníkov na doručenia, ktoré nikdy nebudú úspeš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,))
Eskalácia half_open → open s dvojnásobným časom ochladzovania je detail, ktorý väčšina implementácií prehliada. Ak koncový bod zlyhá počas sondy (stav half_open), nechcete to skúsiť znova o ďalších 5 minút. Koncový bod je stále nefunkčný. Zdvojnásobte čas ochladzovania na 10 minút, potom na 20, s obmedzením na 1 hodinu. Tým sa zabráni tomu, aby sa istič stal periodickým mechanizmom kladiva.
Sledovanie percentilu latencie
Priemery klamú. Koncový bod s priemernou dobou odozvy 200 ms môže reagovať za 50 ms v 95 % prípadov a za 3 000 ms v zvyšných 5 %. Priemer vyzerá dobre. P95 odhaľuje problém, ktorý 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á funkcia usporiadaných množín, ktorá počíta presné percentily. Pre veľké množiny údajov by ste prešli na percentile_disc (ktorá vracia skutočnú pozorovanú hodnotu namiesto interpolácie) alebo použite aproximáciu t-digest. Pre doručovanie webhookov s 24-hodinovým oknom sú presné percentily na nespracovaných údajoch dostatočne rýchle až do približne 100 000 meraní na koncový bod.
get_slow_endpoints Funkciu spúšťam ako plánovanú kontrolu každých 15 minút. Koncové body s P95 nad 2 sekundy sú označené na vyšetrovanie. Koncovým bodom s P95 nad 5 sekúnd sa znižuje prah ističa – je im povolený menší počet po sebe nasledujúcich zlyhaní pred otvorením okruhu, pretože každé neúspešné doručenie viaže pracovné vlákno na celé trvanie časového limitu.
Monitorovanie stavu doručovania webhookov
Tu je monitorovací dotaz, ktorý spúšťam každých päť minút. Vytvorí jednoriadkový súhrn stavu celého doručovacieho kanála:
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 Hodnota je najdôležitejšou metrikou v tomto dotaze. Ak je staršia ako vaše maximálne okno opakovania (súčet všetkých oneskorení odpočítavania), niečo je štrukturálne nesprávne – buď je pracovník zaseknutý, koncový bod je blokovaný alebo front rastie rýchlejšie, ako ho stihnete vyprázdniť.
Upozorňujem na tri podmienky: zvyšujúci sa počet nedoručených listov (koncové body trvalo zlyhávajú a nikto to neskúma), vek čakajúceho frontu presahuje 30 minút (doručenie sa oneskoruje) a Latencia P95 na koncový bod prekračuje prahové hodnoty, ktoré naznačujú zníženú spoľahlivosť doručovaniaTretím je signál včasného varovania – latencia sa zvyšuje ešte pred zlyhaním. Koncový bod, ktorý reagoval 200 ms a začne reagovať o 3 sekundy, čoskoro začne vypršať časový limit.
Rad mŕtvych listov nie je len úložisko
Väčšina tímov implementuje front mŕtvych listov ako tabuľku, kde sa zlyhávajú webhooky. Občas ju kontrolujú počas reakcie na incident. To je plytvanie.
Front mŕtvych listov je vaša najcennejšia ladiaca množina údajov. Každý riadok predstavuje doručenie, ktoré sa váš systém pokúsil doručiť viackrát a vzdal sa ho. Vzor mŕtvych listov vám prezradí veci, ktoré vám metriky úspešnosti nikdy neposkytnú.
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()
Keď kontrolujem mŕtve listy, hľadám tri vzorce.
Zlyhania klastra: 50 mŕtvych písmen pre ten istý koncový bod v tej istej hodine znamená, že koncový bod prestal fungovať a neobnovil sa v rámci okna na opakovanie. Akcia: predĺžte okno na opakovanie alebo implementujte manuálne opätovné zaradenie do frontu.
Vzory stavových kódov: Nárast počtu neplatných hlásení 401/403 znamená, že koncový bod zmenil poverenia a nikto neaktualizoval konfiguráciu webhooku. Nárast počtu hlásení 429 (príliš veľa požiadaviek) znamená, že prekračujete ich limit rýchlosti a je potrebné obmedziť ich používanie.
Postupná akumulácia: 2 – 3 mŕtve listy denne pre jeden koncový bod, rovnomerne rozložené. Toto je najzákernejší vzorec – koncový bod väčšinou funguje, ale má občasné zlyhania, ktoré časom vyčerpávajú opakované pokusy. Oprava sa zvyčajne zvyšuje. max_attempts pre daný konkrétny koncový bod alebo skrátenie časového limitu.
Spustenie pracovníka
Hlavná slučka, ktorá všetko spája:
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()
SIGTERM Obslužný program je nevyhnutný pre čisté ukončenia v kontajnerových prostrediach. Keď Kubernetes odošle SIGTERM, pracovník dokončí svoju aktuálnu dávku, potvrdí transakciu a ukončí sa. Bez neho sa riadky zaseknú v in_flight stav bez pracovníka, ktorý ich spracováva.
Čo tento systém nerobí (a kedy potrebujete viac)
Táto implementácia zvládne až približne 10 000 doručení za minútu na jednej inštancii PostgreSQL s 2 – 3 pracovnými procesmi. Okrem toho sú potrebné tri zmeny.
Najprv nahraďte front PostgreSQL streammi Redis alebo RabbitMQ. SELECT FOR UPDATE SKIP LOCKED Tento vzor vytvára konflikt zápisu v tabuľke frontu pri vysokej priepustnosti. Vyhradený sprostredkovateľ správ to eliminuje.
Po druhé, pridajte obmedzenie rýchlosti pre každý koncový bod. Niektoré prijímajúce koncové body majú obmedzenia rýchlosti (100 požiadaviek za minútu, 1 000 za hodinu). Bez obmedzenia rýchlosti na strane klienta prekročíte ich kvótu a dostanete chybu 429. Implementujte token bucket pre každý koncový bod.
Po tretie, pridajte podpisovanie požiadaviek. Podpisy HMAC-SHA256 v úžitkovej záťaži umožňujú prijímajúcemu koncovému bodu overiť, či webhook pochádza z vášho systému a nebol počas prenosu zmenený. Toto je dôležité pre akýkoľvek systém webhookov, ktorý odosiela finančné údaje.
Systém v tomto článku je základom. Rieši zložité problémy – logiku opakovaných pokusov, prerušenie okruhu, meranie latencie, analýzu nedoručených správ – ktoré každý systém doručovania webhookov potrebuje bez ohľadu na rozsah. Konkrétne komponenty, ktoré pripojíte (sprostredkovateľ správ, obmedzovač rýchlosti, podpisovanie požiadaviek), závisia od vašich požiadaviek na priepustnosť a bezpečnosť.
Najdôležitejšia časť je tá, ktorú väčšina tímov preskakuje: meranie samotného systému doručovania. Ak neviete odpovedať na otázku „aká je latencia doručenia P95 do koncového bodu X za posledných 24 hodín“, pracujete naslepo. Najprv si zostavte inštrumentáciu. Všetko ostatné nasleduje.
Často kladené otázky o doručovaní webhookov
Mal by odosielateľ webhooku sľúbiť doručenie presne raz?
Zvyčajne nie. Odosielateľ by mal zviditeľniť opakované pokusy a poskytnúť stabilné ID udalosti; príjemca by mal nastaviť spracovanie ako idempotentné, aby sa udalosť mohla doručiť viackrát bez duplikovania obchodného výsledku.
Ktoré zlyhania by sa mali zopakovať?
Zlyhania, ktoré vaša zmluva klasifikuje ako prechodné, ako sú napríklad sieťové chyby, časové limity a vybrané odpovede servera, opakujte iba v prípade chybných požiadaviek, zlyhaní overenia alebo iných trvalých chýb bez explicitnej cesty k náprave.
Kedy by sa mal tím presunúť za rámec frontu podporovaného databázou?
Presunúť sa, keď namerané súboje, vek nevybavených objednávok, priepustnosť alebo potreby operačnej obnovy ukážu, že databázový front už nespĺňa zmluvu o doručení. Rozhodnutia o kapacite by sa mali riadiť pozorovanou pracovnou záťažou, nie všeobecným tvrdením o miere požiadaviek.