Ostatnia aktualizacja 24 czerwca 2026 r. Przez Cezar Fikson
Odpowiedź bezpośrednia: Niezawodna usługa dostarczania webhooków lub afiliacyjnych postbacków wymaga trwałej kolejki, ograniczonej liczby ponownych prób z uwzględnieniem drgań, wyłączania obwodów dla każdego punktu końcowego, ochrony przed duplikatami, weryfikacji podpisów i danych telemetrycznych, które pokazują, czy dostarczanie danych zwalnia, zanim zawiedzie.
Ta implementacja celowo utrzymuje model wykonania na niskim poziomie: Python Workers i PostgreSQL dla kolejki, historii dostaw i stanu punktów końcowych. Stanowi praktyczny punkt wyjścia dla rozwiązań SaaS B2B, działań partnerskich w zakresie gier internetowych iGaming oraz każdego produktu, który musi dostarczać zdarzenia wychodzące bez traktowania pojedynczego przekroczenia limitu czasu HTTP jako utraconej konwersji.
Co musi gwarantować umowa dostawy
| Control: | Dlaczego jest to ważne | Kontrola operacyjna |
|---|---|---|
| Trwały zapis zdarzeń | Wydarzenia są zachowywane po ponownym uruchomieniu pracownika. | Każda przyjęta dostawa ma stały identyfikator i status. |
| Podpisany wniosek | Odbiorcy mogą zweryfikować nadawcę i wykryć próbę manipulacji przy ciele. | Używaj podpisu ze znacznikiem czasu i odrzucaj nieaktualne żądania. |
| Ochrona przed duplikacją | Ponawianie prób może skutkować duplikacją konwersji lub aktualizacji. | Wyślij identyfikator zdarzenia i uczyń odbiorcę idempotentnym. |
| Ograniczone ponowne próby | Przejściowe awarie naprawiają się bez przeciążania zdegradowanego punktu końcowego. | Wycofaj się z drżeniem; zatrzymaj się po udokumentowanym limicie prób. |
| Telemetria punktu końcowego | Sama długość kolejki ukrywa powolnych lub zawodnych partnerów. | Śledź wskaźnik sukcesu, najstarsze oczekujące zdarzenie i percentyle opóźnień dla każdego punktu końcowego. |
Bezpieczeństwo i kontrola duplikatów są ważniejsze niż ponowne dostrajanie
Traktuj treść wychodzącą jako poufne dane operacyjne. Podpisz dokładną treść żądania, dołącz znacznik czasu dostarczenia i niezmienny identyfikator zdarzenia, wymieniaj klucze uwierzytelniające i upewnij się, że aplikacja odbierająca bezpiecznie zignoruje powtórzenie tego samego zdarzenia. Pomyślna odpowiedź HTTP nie jest dowodem na to, że zdarzenie biznesowe zostało zastosowane dokładnie raz; decyzję musi podjąć strona odbierająca.
W przypadku przepływu pracy dotyczącego postbacku afiliacyjnego lub iGamingu, nie zapisuj identyfikatorów kliknięć, identyfikatorów konwersji ani pól istotnych dla wypłat w logach, chyba że kontrola dostępu i reguły przechowywania danych są wyraźnie określone. Poniższe odniesienie do Scaleo ma znaczenie jedynie jako przykład opóźnienia postbacku afiliacyjnego; nie zastępuje ono dokumentowania umowy dostawy między Twoimi własnymi usługami.
Przydatne odniesienia dotyczące implementacji: Weryfikacja podpisu webhooka Stripe, PostgreSQL SELECT i SKIP ZABLOKOWANE, Wskazówki AWS dotyczące wykładniczego wycofywania i jittera.
Model danych
Wszystko zaczyna się od dwóch tabel: jednej dla kolejki dostaw i jednej dla pomiarów opóźnień.
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
);
Trzy rzeczy, na które należy zwrócić uwagę w tym schemacie.
Po pierwsze, webhook_queue tabela używa next_attempt_at kolumna zamiast oddzielnego mechanizmu planowania. Pracownik sprawdza wiersze, w których status IN ('pending', 'failed') AND next_attempt_at <= NOW()To jest kolejka opóźniona dla ubogich i działa dobrze do około 10 000 przesyłek na minutę. Powyżej tej liczby należy zainstalować odpowiedniego brokera wiadomości.
Po drugie, endpoint_latency Tabela działa jak bufor pierścieniowy. Okresowo usuwam wiersze starsze niż 24 godziny. Percentyle opóźnienia w endpoint_health są obliczane na podstawie tego ruchomego okna — reprezentują ostatnie zachowania, a nie historyczne średnie.
Po trzecie, endpoint_health Tabela implementuje maszynę stanów wyłącznika. Więcej na ten temat poniżej.
Pracownik dostawy
Główna pętla robocza jest celowo prosta. Złożoność leży w logice ponawiania prób i wyłączniku obwodu, a nie w samej ścieżce dostarczania.
pyton
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 jest krytyczna dla uruchamiania wielu instancji roboczych. Bez SKIP LOCKEDDwóch pracowników blokowałoby się w tym samym wierszu. Dzięki temu każdy pracownik pobiera inną partię oczekujących webhooków. Zapewnia to skalowalność poziomą poprzez proste uruchamianie większej liczby procesów roboczych.
time.monotonic() zadzwoń zamiast time.time() jest celowe. time.time() może cofnąć się podczas regulacji NTP. time.monotonic() nigdy się nie cofa, co ma znaczenie, gdy mierzymy opóźnienia mniejsze niż sekunda.
Logika ponawiania prób z wykładniczym odliczaniem i drżeniem
Gdy dostarczenie danych się nie powiedzie, czas ponownych prób decyduje o tym, czy system odzyska sprawność działania, czy też utworzy się grupa urządzeń, które zaatakuje mający problemy punkt końcowy.
pyton
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)
Dlaczego pełny jitter zamiast jittera dekorrelowanego lub równego jittera? AWS opublikował ostateczną analizę na ten temat. Pełny jitter (losowanie między 0 a limitem wykładniczym) zapewnia najniższy całkowity czas ukończenia dla wszystkich klientów. Równy jitter (losowanie między połową limitu a pełnym limitem) jest bardziej konserwatywny, ale wolniej wyczerpuje listę ponownych prób. W przypadku dostarczania webhooków z wieloma niezależnymi punktami końcowymi, pełny jitter jest właściwym wyborem, ponieważ ponowne próby każdego punktu końcowego są niezależne — nie są one koordynowane.
Wyłącznik obwodu: Przestań niszczyć uszkodzone punkty końcowe
Wzorzec wyłącznika zapobiega marnowaniu zasobów systemu na punkty końcowe, które stale ulegają awariom. Bez niego martwy punkt końcowy gromadzi setki oczekujących ponownych prób, z których każda kończy się 15-sekundowym limitem czasu – marnując moce przerobowe pracowników na dostawy, które nigdy się nie powiodą.
pyton
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,))
Eskalacja half_open → open z podwojonym czasem odnowienia to szczegół, którego brakuje w większości implementacji. Jeśli punkt końcowy ulegnie awarii podczas sondowania (stan half_open), nie chcesz ponawiać próby za kolejne 5 minut. Punkt końcowy nadal jest uszkodzony. Podwój czas odnowienia do 10 minut, następnie do 20, aż do 1 godziny. Zapobiega to przekształcaniu się wyłącznika w mechanizm okresowego generowania awarii.
Śledzenie percentyla opóźnienia
Średnie kłamią. Punkt końcowy ze średnim czasem reakcji 200 ms może odpowiadać w 50 ms w 95% przypadków i 3,000 ms w pozostałych 5%. Średnia wygląda dobrze. P95 ujawnia problem, który dotyczy 1 na 20 dostaw.
pyton
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 to funkcja agregująca o uporządkowanym zbiorze, która oblicza dokładne percentyle. W przypadku dużych zbiorów danych należy przełączyć się na percentile_disc (który zwraca rzeczywistą wartość obserwowaną, a nie interpolację) lub użyć przybliżenia t-digest. W przypadku dostarczania webhooków z 24-godzinnym oknem, dokładne percentyle surowych danych są wystarczająco szybkie do około 100 000 pomiarów na punkt końcowy.
get_slow_endpoints Funkcja ta jest uruchamiana przeze mnie jako zaplanowane sprawdzenie co 15 minut. Punkty końcowe z wartością P95 powyżej 2 sekund są oznaczane do sprawdzenia. Punkty końcowe z wartością P95 powyżej 5 sekund mają obniżony próg zadziałania wyłącznika obwodu — dopuszcza się mniej kolejnych awarii przed otwarciem obwodu, ponieważ każda nieudana próba dostarczenia danych blokuje wątek roboczy na cały czas trwania limitu czasu.
Monitorowanie stanu dostarczania webhooków
Oto zapytanie monitorujące, które uruchamiam co pięć minut. Generuje ono jednowierszowe podsumowanie stanu całego potoku dostaw:
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 Wartość to najważniejsza metryka w tym zapytaniu. Jeśli jest starsza niż maksymalne okno ponownych prób (suma wszystkich opóźnień), coś jest nie tak ze strukturą — albo proces roboczy się zawiesił, punkt końcowy jest zablokowany, albo kolejka rośnie szybciej, niż jesteś w stanie ją opróżnić.
Ostrzegam o trzech sytuacjach: zwiększająca się liczba martwych wiadomości (punkty końcowe stale ulegają awarii i nikt nie zajmuje się badaniem problemu), wiek oczekującej kolejki przekraczający 30 minut (dostarczanie jest opóźnione) oraz opóźnienie P95 dla każdego punktu końcowego przekraczające progi wskazujące na obniżoną niezawodność dostarczaniaTrzecim sygnałem jest wczesny sygnał ostrzegawczy – opóźnienie wzrasta przed wystąpieniem awarii. Punkt końcowy, który odpowiadał w ciągu 200 ms i zaczyna odpowiadać w ciągu 3 sekund, wkrótce zacznie przekraczać limit czasu.
Kolejka martwych listów to nie tylko magazyn
Większość zespołów implementuje kolejkę martwych wiadomości jako tabelę, do której trafiają niedziałające webhooki. Sprawdzają ją od czasu do czasu podczas reagowania na incydenty. To marnotrawstwo.
Kolejka martwych listów to Twój najcenniejszy zbiór danych do debugowania. Każdy wiersz reprezentuje doręczenie, które Twój system próbował wielokrotnie i z którego zrezygnował. Wzór martwych listów ujawnia rzeczy, których metryki sukcesu nigdy nie ujawnią.
pyton
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()
Przeglądając martwe listy, szukam trzech wzorców.
Awarie klastra: 50 martwych wiadomości dla tego samego punktu końcowego w ciągu tej samej godziny oznacza, że punkt końcowy uległ awarii i nie odzyskał sprawności w ciągu okna ponownych prób. Działanie: wydłuż okno ponownych prób lub wprowadź ręczne ponowne kolejkowanie.
Wzory kodów statusu: Skok liczby martwych wiadomości 401/403 oznacza, że punkt końcowy dokonał rotacji danych uwierzytelniających i nikt nie zaktualizował konfiguracji webhooka. Skok liczby 429 (zbyt wiele żądań) oznacza, że przekraczasz limit przepustowości i musisz ograniczyć przepustowość.
Stopniowa akumulacja: 2-3 martwe listy dziennie dla pojedynczego punktu końcowego, równomiernie rozłożone. To najbardziej podstępny schemat — punkt końcowy działa w większości przypadków, ale zdarzają się sporadyczne awarie, które z czasem wyczerpują możliwości ponownych prób. Naprawa zazwyczaj rośnie. max_attempts dla danego punktu końcowego lub skrócenia limitu czasu.
Uruchamianie pracownika
Główna pętla, która wszystko łączy:
pyton
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 Obsługa jest niezbędna do czystego zamykania w środowiskach konteneryzowanych. Gdy Kubernetes wysyła sygnał SIGTERM, pracownik kończy bieżącą partię, zatwierdza transakcję i wychodzi. Bez tego wiersze utkną w… in_flight status bez pracownika, który by je przetwarzał.
Czego ten system nie robi (i kiedy potrzebujesz czegoś więcej)
Ta implementacja obsługuje do około 10 000 przesyłek na minutę na jednej instancji PostgreSQL z 2-3 procesami roboczymi. Powyżej tej liczby wymagane są trzy zmiany.
Najpierw zastąp kolejkę PostgreSQL Redis Streams lub RabbitMQ. SELECT FOR UPDATE SKIP LOCKED Wzorzec powoduje konflikt zapisu w tabeli kolejki przy dużej przepustowości. Dedykowany broker komunikatów eliminuje to zjawisko.
Po drugie, dodaj limity przepustowości dla każdego punktu końcowego. Niektóre punkty odbiorcze mają limity przepustowości (100 żądań na minutę, 1,000 na godzinę). Bez limitów przepustowości po stronie klienta przekroczysz ich limit i otrzymasz 429 żądań. Wdróż kontener tokenów dla każdego punktu końcowego.
Po trzecie, dodaj podpisywanie żądań. Podpisy HMAC-SHA256 w treści pozwalają punktowi końcowemu odbiorcy zweryfikować, czy webhook pochodzi z Twojego systemu i nie został zmodyfikowany w trakcie transmisji. To podstawa dla każdego systemu webhook, który wysyła dane finansowe.
System opisany w tym artykule stanowi fundament. Rozwiązuje on trudne problemy – logikę ponawiania prób, wyłączanie obwodów, pomiar opóźnień, analizę martwych komunikatów – których potrzebuje każdy system dostarczania webhooków, niezależnie od skali. Konkretne komponenty, które należy dodać (broker komunikatów, ogranicznik prędkości, podpisywanie żądań), zależą od przepustowości i wymagań bezpieczeństwa.
Najważniejsza jest ta część, którą większość zespołów pomija: pomiar samego systemu dostarczania. Jeśli nie potrafisz odpowiedzieć na pytanie „Jakie jest opóźnienie w dostarczaniu P95 do Endpoint X w ciągu ostatnich 24 godzin”, działasz w ciemno. Najpierw zbuduj instrumentację. Wszystko inne nastąpi później.
Często zadawane pytania dotyczące dostarczania webhooków
Czy nadawca webhooku powinien obiecać dostarczenie wiadomości dokładnie raz?
Zazwyczaj nie. Nadawca powinien uwidocznić ponowne próby i podać stabilny identyfikator zdarzenia; odbiorca powinien ustawić przetwarzanie idempotentne, aby zdarzenie mogło zostać dostarczone więcej niż raz bez duplikowania wyniku biznesowego.
Które błędy należy powtórzyć?
Ponawiaj tylko błędy sklasyfikowane w umowie jako przejściowe, takie jak błędy sieciowe, przekroczenia limitu czasu i odpowiedzi wybranych serwerów. Nie ponawiaj wielokrotnie błędnych żądań, błędów uwierzytelniania ani innych trwałych błędów bez wyraźnej ścieżki naprawczej.
Kiedy zespół powinien wyjść poza kolejkę opartą na bazie danych?
Przenieś, gdy pomiary rywalizacji, wieku zaległości, przepustowości lub potrzeb odzyskiwania operacyjnego wskazują, że kolejka bazy danych nie spełnia już warunków umowy dostawy. Decyzje dotyczące pojemności powinny być podejmowane na podstawie obserwowanego obciążenia, a nie na podstawie ogólnego zapotrzebowania na żądanie.