Construindo um Sistema de Entrega de Webhooks: Postbacks Confiáveis ​​em 2026

Um guia prático para 2026 sobre webhooks duráveis ​​e entrega de posts de afiliados com filas, novas tentativas com jitter, disjuntores, idempotência, assinaturas e telemetria de latência por endpoint.
Sistema de entrega de webhooks - Construindo um sistema de entrega de webhooks: Postbacks confiáveis ​​em 2026

Atualizado pela última vez em 24 de junho de 2026 por César Fikson

Resposta direta: Um serviço confiável de entrega via webhook ou postback de afiliados precisa de uma fila robusta, número limitado de tentativas com jitter, circuito de interrupção por ponto de extremidade, proteção contra duplicatas, verificação de assinatura e telemetria que mostre se a entrega está ficando mais lenta antes que falhe.

Esta implementação mantém o modelo de execução intencionalmente pequeno: workers em Python e PostgreSQL para fila, histórico de entregas e integridade dos endpoints. É um ponto de partida prático para SaaS B2B, operações de afiliados de iGaming e qualquer produto que precise entregar eventos de saída sem considerar um único timeout HTTP como uma conversão perdida.

O que o contrato de entrega deve garantir

Controlar Por que é importante verificação operacional
Registro de eventos duradouros Os eventos persistem mesmo após reinicializações do processo. Cada entrega aceita possui um ID e um status estáveis.
Solicitação assinada Os receptores podem verificar o remetente e detectar adulteração do corpo da embalagem. Use uma assinatura com registro de data e hora e rejeite solicitações obsoletas.
Proteção contra duplicados Caso contrário, novas tentativas podem criar conversões ou atualizações duplicadas. Envie um ID de evento e torne o receptor idempotente.
Tentativas limitadas Falhas transitórias se recuperam sem sobrecarregar um ponto final degradado. Reduza a intensidade com moderação; pare após um limite de tentativas documentado.
Telemetria de ponto final A profundidade da fila, por si só, oculta parceiros lentos ou com falhas. Acompanhe a taxa de sucesso, o evento pendente mais antigo e os percentis de latência por ponto de extremidade.

Segurança e controles de duplicados vêm antes da nova tentativa de ajuste.

Trate o corpo da requisição de saída como dados operacionais sensíveis. Assine o corpo da requisição exatamente como está, inclua um registro de data e hora de entrega e um ID de evento imutável, alterne os segredos de assinatura e assegure-se de que o aplicativo receptor ignore com segurança a repetição do mesmo evento. Uma resposta HTTP bem-sucedida não comprova que um evento de negócio foi aplicado exatamente uma vez; essa decisão cabe ao lado receptor.

Para um fluxo de trabalho de postback de afiliados ou iGaming, mantenha os IDs de cliques, IDs de conversão e campos relevantes para pagamentos fora dos logs, a menos que os controles de acesso e as regras de retenção sejam explícitos. A referência da Scaleo existente abaixo é relevante apenas como um exemplo de latência de postback de afiliados; ela não substitui a documentação do contrato de entrega entre seus próprios serviços.

Referências úteis para implementação: Verificação de assinatura via webhook do Stripe, PostgreSQL SELECT e SKIP BLOQUEADO e Orientações da AWS sobre recuo exponencial e jitter..

O Modelo de Dados

Tudo começa com duas tabelas: uma para a fila de entrega e outra para as medições de latência.

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

Três pontos a serem observados sobre este esquema.

Em primeiro lugar, o webhook_queue tabela usa um next_attempt_at coluna em vez de um mecanismo de agendamento separado. O trabalhador consulta as linhas onde status IN ('pending', 'failed') AND next_attempt_at <= NOW()Esta é uma solução improvisada para filas de atraso e funciona bem até cerca de 10,000 entregas por minuto. Acima disso, utilize um broker de mensagens adequado.

Em segundo lugar, o endpoint_latency A tabela funciona como um buffer circular. Periodicamente, excluo linhas com mais de 24 horas. Os percentis de latência em endpoint_health são calculados a partir dessa janela móvel — eles representam o comportamento recente, não as médias históricas.

Em terceiro lugar, o endpoint_health A tabela implementa a máquina de estados do disjuntor. Mais detalhes abaixo.

O entregador

O loop principal do worker é propositalmente simples. A complexidade reside na lógica de repetição e no disjuntor, não no próprio caminho de entrega.

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

O SELECT FOR UPDATE SKIP LOCKED Essa cláusula é crucial para executar várias instâncias de trabalho. Sem ela, sem SKIP LOCKED, dois workers ficariam bloqueados na mesma linha. Com isso, cada worker obtém um lote diferente de webhooks pendentes. Isso permite escalabilidade horizontal simplesmente iniciando mais processos worker.

O time.monotonic() ligar em vez de time.time() É intencional. time.time() Pode haver retrocessos durante ajustes de NTP. time.monotonic() nunca retrocede, o que é importante quando se está medindo latência inferior a um segundo.

Lógica de repetição com recuo exponencial e jitter

Quando uma entrega falha, o tempo de nova tentativa determina se o seu sistema se recupera de forma adequada ou se cria uma sobrecarga que afeta um ponto de extremidade com dificuldades.

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)

Por que usar jitter completo em vez de jitter decorrelacionado ou jitter igual? A AWS publicou a análise definitiva sobre isso. O jitter completo (aleatorizando entre 0 e o limite exponencial) produz o menor tempo total de conclusão para todos os clientes. O jitter igual (aleatorizando entre metade do limite e o limite total) é mais conservador, mas mais lento para reduzir o backlog de tentativas. Para entrega de webhooks com muitos endpoints independentes, o jitter completo é a escolha certa, pois as tentativas de cada endpoint são independentes — não há coordenação entre elas.

Disjuntor: Pare de martelar pontos finais quebrados

O padrão de disjuntor impede que seu sistema desperdice recursos em endpoints que falham constantemente. Sem ele, um endpoint inativo acumula centenas de tentativas pendentes que expiram em 15 segundos cada, consumindo a capacidade de processamento em entregas que nunca serão concluídas.

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

A escalação de estado "semiaberto" para "aberto" com o dobro do tempo de espera é o detalhe que a maioria das implementações ignora. Se um endpoint falhar durante a sondagem (estado "semiaberto"), você não vai querer tentar novamente em 5 minutos. O endpoint ainda estará com problemas. Dobre o tempo de espera para 10 minutos, depois para 20, com um limite máximo de 1 hora. Isso impede que o disjuntor se torne um mecanismo de sobrecarga periódico.

Rastreamento de Percentil de Latência

As médias enganam. Um endpoint com um tempo médio de resposta de 200 ms pode responder em 50 ms em 95% das vezes e em 3,000 ms nos outros 5%. A média parece boa. O P95 revela um problema que afeta 1 em cada 20 entregas.

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 percentile_cont é uma função agregada de conjunto ordenado que calcula percentis exatos. Para conjuntos de dados grandes, você pode usar percentile_disc (que retorna um valor observado real em vez de interpolar) ou usar uma aproximação de resumo t. Para entrega via webhook com janela de 24 horas, percentis exatos nos dados brutos são suficientemente rápidos para até cerca de 100,000 medições por endpoint.

O get_slow_endpoints A função é o que eu executo como uma verificação agendada a cada 15 minutos. Os endpoints com P95 acima de 2 segundos são sinalizados para investigação. Os endpoints com P95 acima de 5 segundos têm seu limite de disjuntor reduzido — eles são permitidos a ocorrer menos falhas consecutivas antes que o circuito seja aberto, porque cada entrega com falha ocupa um thread de trabalho durante toda a duração do tempo limite.

Monitoramento da integridade da entrega do webhook

Aqui está a consulta de monitoramento que executo a cada cinco minutos. Ela gera um resumo de integridade em uma única linha de todo o pipeline de entrega:

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;

O oldest_pending O valor é a métrica mais importante nesta consulta. Se for mais antigo que sua janela máxima de repetição (soma de todos os atrasos de espera), algo está estruturalmente errado — ou o worker está travado, o endpoint está bloqueado ou a fila está crescendo mais rápido do que você consegue esvaziá-la.

Alerto para três condições: aumento do número de mensagens não entregues (os endpoints estão apresentando falhas permanentes e ninguém está investigando), tempo de espera na fila superior a 30 minutos (a entrega está atrasada) e Latência P95 por ponto final excedendo os limites que indicam degradação da confiabilidade da entrega. O terceiro é o sinal de alerta precoce — a latência aumenta antes que as falhas ocorram. Um endpoint que respondia em 200 ms e passa a responder em 3 segundos está prestes a atingir o tempo limite.

A fila de mensagens não entregues não é apenas um espaço de armazenamento.

A maioria das equipes implementa uma fila de mensagens não entregues como uma tabela onde os webhooks com falha são armazenados. Elas a verificam ocasionalmente durante a resposta a incidentes. Isso é um desperdício.

A fila de mensagens não entregues é o seu conjunto de dados de depuração mais valioso. Cada linha representa uma entrega que o seu sistema tentou várias vezes e desistiu. O padrão das mensagens não entregues revela informações que as métricas de sucesso jamais revelarão.

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

Ao analisar cartas não recebidas, procuro por três padrões.

Falhas no cluster: 50 mensagens não entregues para o mesmo endpoint na mesma hora significa que o endpoint ficou inativo e não se recuperou dentro do período de novas tentativas. Ação: estender o período de novas tentativas ou implementar o reenvio manual de mensagens.

Padrões de código de status: Um pico de mensagens não entregues (401/403) significa que o endpoint rotacionou as credenciais e ninguém atualizou a configuração do webhook. Um pico de mensagens 429 (Muitas Requisições) significa que você está excedendo o limite de requisições e precisa reduzir a taxa de envios.

Acumulação gradual: De duas a três mensagens não entregues por dia para um único endpoint, distribuídas uniformemente. Este é o padrão mais traiçoeiro: o endpoint funciona na maior parte do tempo, mas apresenta falhas intermitentes que esgotam as tentativas ao longo do tempo. A solução geralmente é aumentar o número de tentativas. max_attempts para esse ponto de extremidade específico ou reduzindo o tempo limite.

Executando o trabalhador

O principal elemento que une tudo:

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

O SIGTERM O manipulador é essencial para desligamentos limpos em ambientes conteinerizados. Quando o Kubernetes envia um sinal SIGTERM, o worker finaliza seu lote atual, confirma a transação e encerra. Sem isso, você terá linhas presas em um loop infinito. in_flight status sem nenhum trabalhador processando-os.

O que este sistema não faz (e quando você precisa de mais)

Esta implementação processa até cerca de 10,000 entregas por minuto em uma única instância do PostgreSQL com 2 a 3 processos de trabalho. Acima desse limite, são necessárias três alterações.

Primeiro, substitua a fila do PostgreSQL pelo Redis Streams ou RabbitMQ. SELECT FOR UPDATE SKIP LOCKED O padrão gera contenção de escrita na tabela de filas em situações de alta taxa de transferência. Um agente de mensagens dedicado elimina esse problema.

Em segundo lugar, adicione limitação de taxa por endpoint. Alguns endpoints de recebimento têm limites de taxa (100 requisições por minuto, 1,000 por hora). Sem limitação de taxa no lado do cliente, você ultrapassará a cota deles e receberá um erro 429. Implemente um bucket de tokens por endpoint.

Em terceiro lugar, adicione a assinatura de requisição. As assinaturas HMAC-SHA256 no payload permitem que o endpoint receptor verifique se o webhook veio do seu sistema e não foi adulterado durante a transmissão. Isso é essencial para qualquer sistema de webhook que envie dados financeiros.

O sistema descrito neste artigo é a base. Ele lida com os problemas complexos — lógica de repetição, disjuntor, medição de latência, análise de mensagens não entregues — que todo sistema de entrega de webhooks precisa, independentemente da escala. Os componentes específicos que você adiciona (corretor de mensagens, limitador de taxa, assinatura de requisições) dependem dos seus requisitos de throughput e segurança.

A parte mais importante é aquela que a maioria das equipes ignora: medir o próprio sistema de entrega. Se você não consegue responder à pergunta “qual é a latência de entrega P95 para o Endpoint X nas últimas 24 horas?”, você está operando às cegas. Construa a instrumentação primeiro. Todo o resto vem depois.

Perguntas frequentes sobre a entrega de webhooks

O remetente de um webhook deve prometer entrega exatamente uma única vez?

Normalmente não. O remetente deve tornar as novas tentativas visíveis e fornecer um ID de evento estável; o destinatário deve tornar o processamento idempotente para que um evento possa ser entregue mais de uma vez sem duplicar o resultado comercial.

Quais falhas devem ser repetidas?

Tente novamente apenas falhas que seu contrato classifica como transitórias, como erros de rede, tempos limite e respostas selecionadas do servidor. Não tente novamente solicitações malformadas, falhas de autenticação ou outros erros permanentes repetidamente sem um caminho de correção explícito.

Quando uma equipe deve ir além de uma fila baseada em banco de dados?

Mova o banco de dados quando a contenção medida, a idade do backlog, a taxa de transferência ou as necessidades de recuperação operacional demonstrarem que a fila do banco de dados não atende mais ao contrato de entrega. As decisões sobre capacidade devem ser baseadas na carga de trabalho observada, e não em uma estimativa genérica da taxa de requisições.

Artigo Anterior

Melhor software de rastreamento de afiliados de iGaming em 2026

Próximo Artigo

Inteligência Artificial em Jogos Online: Casos de Uso do ChatGPT para Cassinos

César Fikson
Autor:

César Fikson

Sou Analista de Dados de iGaming, especializado em examinar e interpretar dados relacionados a plataformas de jogos online e atividades de apostas, bem como tendências de mercado. Analiso o comportamento do jogador, o desempenho do jogo e as tendências de receita para otimizar experiências de jogo e estratégias de negócios.

Solicite uma demonstração
PASSO 1 DE 3
Obrigado — você está na fila.
Um engenheiro de soluções da NowG entrará em contato em até um dia útil para agendar sua demonstração.
Índice