Ú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 degradadaLa 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.