Creación de un sistema de entrega de webhooks: Postbacks fiables en 2026

Una guía práctica para 2026 sobre la entrega duradera de webhooks y postbacks de afiliados con colas, reintentos con fluctuación, disyuntores, idempotencia, firmas y telemetría de latencia por punto final.
Sistema de entrega de webhooks: Creación de un sistema de entrega de webhooks: Postbacks fiables en 2026

Última actualización el 24 de junio de 2026 por César Fikson

Respuesta directa: Un servicio de entrega fiable mediante webhooks o postbacks de afiliados necesita una cola duradera, reintentos limitados con fluctuación, interrupción de circuito por punto final, protección contra duplicados, verificación de firma y telemetría que muestre si la entrega se está ralentizando antes de que falle.

Esta implementación mantiene el modelo de ejecución intencionadamente reducido: procesos Python y PostgreSQL para la cola, el historial de entregas y el estado de los puntos finales. Es un punto de partida práctico para SaaS B2B, operaciones de afiliados de iGaming y cualquier producto que deba enviar eventos salientes sin considerar un único tiempo de espera HTTP como una conversión perdida.

Lo que debe garantizar el contrato de entrega

Control Por qué importa Verificación operativa
Registro de eventos duradero Los eventos sobreviven a los reinicios de los trabajadores. Cada entrega aceptada tiene una identificación y un estado estables.
Solicitud firmada Los receptores pueden verificar al remitente y detectar cualquier manipulación del cuerpo. Utilice una firma con marca de tiempo y rechace las solicitudes obsoletas.
Protección contra duplicados De lo contrario, los reintentos pueden generar conversiones o actualizaciones duplicadas. Envía un ID de evento y haz que el receptor sea idempotente.
Reintentos limitados Los fallos transitorios se recuperan sin sobrecargar un punto final degradado. Retroceda con fluctuaciones; deténgase después de un límite de intentos documentado.
Telemetría de puntos finales La profundidad de la cola por sí sola oculta socios lentos o que están fallando. Realizar un seguimiento de la tasa de éxito, el evento pendiente más antiguo y los percentiles de latencia por punto final.

La seguridad y los controles de duplicados tienen prioridad sobre la optimización de reintentos.

Trate el cuerpo de la solicitud saliente como datos operativos confidenciales. Firme el cuerpo exacto de la solicitud, incluya una marca de tiempo de entrega y un ID de evento inmutable, rote las claves de firma y asegúrese de que la aplicación receptora ignore de forma segura la repetición del mismo evento. Una respuesta HTTP exitosa no prueba que un evento de negocio se haya aplicado exactamente una vez; la parte receptora debe tomar esa decisión.

Para flujos de trabajo de postback de afiliados o iGaming, evite incluir en los registros los ID de clic, los ID de conversión y los campos relevantes para el pago, a menos que se especifiquen explícitamente los controles de acceso y las reglas de retención. La referencia a Scaleo que se muestra a continuación solo sirve como ejemplo de latencia de postback de afiliados; no sustituye la documentación del contrato de entrega entre sus propios servicios.

Referencias útiles para la implementación: Verificación de firma mediante webhook de Stripe, PostgreSQL SELECT y SKIP LOCKED, y Guía de AWS sobre retroceso exponencial y fluctuación.

El modelo de datos

Todo comienza con dos tablas: una para la cola de entrega y otra para las mediciones de latencia.

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

Hay tres cosas que cabe destacar sobre este esquema.

En primer lugar, la webhook_queue una tabla utiliza una next_attempt_at columna en lugar de un mecanismo de programación separado. El trabajador sondea las filas donde status IN ('pending', 'failed') AND next_attempt_at <= NOW()Esta es una cola de espera rudimentaria que funciona bien hasta aproximadamente 10 000 entregas por minuto. Para superar ese límite, conviene utilizar un intermediario de mensajes adecuado.

En segundo lugar, la endpoint_latency La tabla actúa como un búfer circular. Periódicamente purgo las filas con más de 24 horas de antigüedad. Los percentiles de latencia en endpoint_health Se calculan a partir de esta ventana móvil; representan el comportamiento reciente, no los promedios históricos.

En tercer lugar, la endpoint_health La tabla implementa la máquina de estados del interruptor automático. Más información a continuación.

El repartidor

El bucle principal del proceso es deliberadamente simple. La complejidad reside en la lógica de reintento y el disyuntor, no en la ruta de entrega en sí.

pitón

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

El SELECT FOR UPDATE SKIP LOCKED Esta cláusula es fundamental para ejecutar múltiples instancias de trabajador. Sin ella, SKIP LOCKEDDos trabajadores se bloquearían en la misma fila. De esta forma, cada trabajador toma un lote diferente de webhooks pendientes. Esto permite una escalabilidad horizontal simplemente iniciando más procesos de trabajador.

El time.monotonic() llamar en lugar de time.time() es intencional. time.time() puede retroceder durante los ajustes de NTP. time.monotonic() Nunca retrocede, lo cual es importante cuando se mide una latencia inferior a un segundo.

Lógica de reintento con retroceso exponencial y fluctuación

Cuando falla una entrega, el tiempo de reintento determina si su sistema se recupera correctamente o si crea una avalancha de solicitudes que sobrecarga un punto final con problemas.

pitón

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 qué usar jitter completo en lugar de jitter descorrelacionado o jitter igual? AWS publicó el análisis definitivo al respecto. El jitter completo (aleatorización entre 0 y el límite exponencial) produce el menor tiempo total de finalización en todos los clientes. El jitter igual (aleatorización entre la mitad del límite y el límite completo) es más conservador, pero tarda más en vaciar la cola de reintentos. Para la entrega de webhooks con muchos endpoints independientes, el jitter completo es la opción correcta porque los reintentos de cada endpoint son independientes; no se requiere coordinación entre ellos.

Disyuntor: Deje de dañar los puntos finales rotos

El patrón de disyuntor evita que su sistema desperdicie recursos en puntos finales que fallan constantemente. Sin él, un punto final inactivo acumula cientos de reintentos pendientes que expiran a los 15 segundos cada uno, consumiendo la capacidad de procesamiento en entregas que nunca se completarán.

pitón

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

La escalada de apertura parcial a apertura completa con un tiempo de espera duplicado es un detalle que la mayoría de las implementaciones pasan por alto. Si un punto final falla durante la comprobación (estado de apertura parcial), no conviene reintentarlo en otros 5 minutos. El punto final sigue sin funcionar. Duplica el tiempo de espera a 10 minutos, luego a 20, con un límite máximo de 1 hora. Esto evita que el interruptor de circuito se convierta en un mecanismo de sobrecarga periódica.

Seguimiento del percentil de latencia

Los promedios pueden ser engañosos. Un dispositivo con un tiempo de respuesta promedio de 200 ms podría responder en 50 ms el 95 % de las veces y en 3,000 ms el 5 % restante. El promedio parece correcto. Sin embargo, el percentil 95 revela un problema que afecta a 1 de cada 20 entregas.

pitón

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 es una función de agregación de conjuntos ordenados que calcula percentiles exactos. Para conjuntos de datos grandes, cambiarías a percentile_disc (que devuelve un valor observado real en lugar de interpolar) o utilice una aproximación t-digest. Para la entrega de webhooks con una ventana de 24 horas, los percentiles exactos en los datos sin procesar son suficientemente rápidos hasta aproximadamente 100 000 mediciones por punto final.

El get_slow_endpoints Esta función se ejecuta como una comprobación programada cada 15 minutos. Los puntos finales con un P95 superior a 2 segundos se marcan para su investigación. A los puntos finales con un P95 superior a 5 segundos se les reduce el umbral del disyuntor: se les permiten menos fallos consecutivos antes de que se abra el circuito, ya que cada entrega fallida ocupa un hilo de trabajo durante todo el tiempo de espera.

Supervisión del estado de entrega del webhook

Esta es la consulta de monitorización que ejecuto cada cinco minutos. Genera un resumen del estado de toda la cadena de entrega en una sola fila:

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;

El oldest_pending El valor es la métrica más importante en esta consulta. Si es anterior a la ventana de reintento máxima (suma de todos los retrasos de retroceso), algo falla estructuralmente: o bien el trabajador está bloqueado, el punto final está en un agujero negro, o la cola crece más rápido de lo que se puede vaciar.

Emito alertas en tres condiciones: aumento del número de cartas fallidas (los puntos finales están fallando permanentemente y nadie está investigando), antigüedad de la cola pendiente superior a 30 minutos (la entrega se está retrasando) y Latencia P95 por punto final que supera los umbrales que indican una fiabilidad de entrega degradada La tercera es la señal de alerta temprana: la latencia aumenta antes de que se produzcan fallos. Un punto final que respondía en 200 ms y empieza a responder en 3 segundos está a punto de agotar el tiempo de espera.

La cola de cartas muertas no es solo almacenamiento

La mayoría de los equipos implementan una cola de mensajes fallidos como una tabla donde los webhooks que fallan terminan su ejecución. La revisan ocasionalmente durante la respuesta a incidentes. Esto es un desperdicio.

La cola de mensajes no entregados es el conjunto de datos más valioso para la depuración. Cada fila representa una entrega que el sistema intentó varias veces y finalmente desistió. El patrón de mensajes no entregados revela información que las métricas de éxito jamás mostrarán.

pitón

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

Cuando reviso cartas sin entregar, busco tres patrones.

Fallos del clúster: 50 mensajes fallidos para el mismo punto final en la misma hora significa que el punto final se cayó y no se recuperó dentro del plazo de reintento. Acción: ampliar el plazo de reintento o implementar una nueva cola manual.

Patrones de códigos de estado: Un aumento repentino en los errores 401/403 indica que el punto final cambió sus credenciales y nadie actualizó la configuración del webhook. Un aumento repentino en los errores 429 (Demasiadas solicitudes) significa que está excediendo su límite de solicitudes y necesita limitarlo.

Acumulación gradual: De 2 a 3 cartas muertas por día para un solo punto final, distribuidas uniformemente. Este es el patrón más engañoso: el punto final funciona en su mayor parte, pero tiene fallas intermitentes que agotan los reintentos con el tiempo. La solución suele ser aumentar max_attempts para ese punto final específico o reduciendo el tiempo de espera.

Ejecutando el trabajador

El bucle principal que lo une todo:

pitón

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

El SIGTERM El manejador es esencial para apagados limpios en entornos de contenedores. Cuando Kubernetes envía SIGTERM, el trabajador finaliza su lote actual, confirma la transacción y sale. Sin esto, se producen filas atascadas en in_flight estado sin ningún trabajador procesándolos.

Lo que este sistema no hace (y cuándo necesitas más)

Esta implementación admite hasta aproximadamente 10 000 entregas por minuto en una única instancia de PostgreSQL con 2 o 3 procesos de trabajo. Para superar este límite, se requieren tres cambios.

Primero, reemplace la cola de PostgreSQL con Redis Streams o RabbitMQ. SELECT FOR UPDATE SKIP LOCKED Este patrón genera contención de escritura en la tabla de colas a alto rendimiento. Un agente de mensajes dedicado elimina este problema.

En segundo lugar, implementa la limitación de velocidad por punto final. Algunos puntos finales receptores tienen límites de velocidad (100 solicitudes por minuto, 1,000 por hora). Sin la limitación de velocidad del lado del cliente, superarás su cuota y recibirás un error 429. Implementa un depósito de tokens por punto final.

En tercer lugar, agregue la firma de la solicitud. Las firmas HMAC-SHA256 en la carga útil permiten que el receptor verifique que el webhook proviene de su sistema y no fue manipulado durante la transmisión. Esto es fundamental para cualquier sistema de webhook que envíe datos financieros.

El sistema descrito en este artículo es la base. Gestiona los problemas complejos —lógica de reintentos, interrupción de circuitos, medición de latencia, análisis de mensajes fallidos— que todo sistema de entrega de webhooks necesita, independientemente de su escala. Los componentes específicos que se añadan (intermediario de mensajes, limitador de velocidad, firma de solicitudes) dependerán del rendimiento y los requisitos de seguridad.

Lo más importante es lo que la mayoría de los equipos omite: medir el sistema de entrega en sí. Si no puedes responder a la pregunta "¿cuál es la latencia de entrega P95 al punto final X en las últimas 24 horas?", estás trabajando a ciegas. Primero, implementa la instrumentación. Todo lo demás vendrá después.

Preguntas frecuentes sobre la entrega de webhooks

¿Debe un remitente de webhook prometer una entrega única?

Normalmente no. El remitente debe hacer visibles los reintentos y proporcionar un ID de evento estable; el receptor debe garantizar que el procesamiento sea idempotente para que un evento pueda entregarse más de una vez sin duplicar el resultado de negocio.

¿Qué fallos deberían repetirse?

Reintente únicamente los fallos que su contrato clasifique como transitorios, como errores de red, tiempos de espera agotados y respuestas de servidores específicos. No reintente repetidamente solicitudes con formato incorrecto, fallos de autenticación u otros errores permanentes sin una solución explícita.

¿Cuándo debería un equipo ir más allá de una cola basada en una base de datos?

Mueva la cola cuando la contención medida, la antigüedad de la cola de espera, el rendimiento o las necesidades de recuperación operativa indiquen que la cola de la base de datos ya no cumple con el contrato de entrega. Las decisiones sobre la capacidad deben basarse en la carga de trabajo observada, no en una tasa de solicitudes genérica.

Artículo anterior

El mejor software de seguimiento de afiliados de iGaming en 2026

Siguiente artículo

Inteligencia artificial en el iGaming: Casos de uso de ChatGPT para casinos

César Fikson
Escrito por

César Fikson

Soy analista de datos de iGaming y me especializo en examinar e interpretar datos relacionados con plataformas de juegos en línea y actividades de apuestas, así como las tendencias del mercado. Analizo el comportamiento de los jugadores, el rendimiento de los juegos y las tendencias de ingresos para optimizar las experiencias de juego y las estrategias comerciales.

Solicita una demo
PASO 1 DE 3
Gracias, estás en la cola.
Un ingeniero de soluciones de NowG se pondrá en contacto con usted en el plazo de un día hábil para programar su visita guiada.
Home