StackPractices
intermediate Por Mathias Paulenko

Patrón Message Deduplication

Prevenir procesamiento duplicado rastreando IDs de mensaje con claves de idempotencia. Los consumidores verifican un almacen antes de procesar para saltar mensajes ya manejados.

Descripción General

Los brokers de mensajes garantizan entrega at-least-once, lo que significa que los consumidores pueden recibir el mismo mensaje mas de una vez. Reintentos de red, caidas del consumidor durante el procesamiento y reentregas del broker causan duplicados. Sin deduplicacion, un pago puede procesarse dos veces o una notificacion enviarse multiples veces.

El patron Message Deduplication rastrea IDs de mensaje en un almacen. Antes de procesar, el consumidor verifica si el ID ya fue manejado. Si si, salta el procesamiento. Si no, procesa el mensaje y registra el ID.

Cuándo Usar

  • For alternatives, see Idempotent Consumer Pattern.

  • Tu broker de mensajes proporciona entrega at-least-once (la mayoria: SQS, RabbitMQ, Kafka)

  • Procesar un mensaje dos veces causa efectos secundarios (pagos, emails, cambios de inventario)

  • Necesitas semantica exactly-once sin un broker que la soporte nativamente

  • Los consumidores caen y los mensajes se reentregan

  • Procesas webhooks de servicios de terceros que pueden reintentar en fallos de red

  • Consumes de multiples colas y necesitas deduplicacion cross-queue

Cuándo Evitar

  • Tu broker proporciona exactly-once nativamente. Las transacciones de Kafka y SQS FIFO con dedup basada en contenido ya manejan esto. Anadir dedup a nivel aplicacion es redundante.
  • Los mensajes son idempotentes por naturaleza. Si procesar dos veces no tiene efectos secundarios (ej., actualizar un timestamp de ultima vista), dedup anade overhead sin valor.
  • El throughput es critico y cada milisegundo cuenta. Dedup con Redis anade ~0.5ms por mensaje. Para sistemas de ultra-baja-latencia, confia en las garantias del broker.
  • No puedes permitirte una dependencia de almacen de dedup. Si el downtime de Redis es inaceptable y no puedes caer a procesamiento idempotente, reconsidera la arquitectura.
  • Los mensajes no tienen ID unico natural. Generar hashes de contenido para cada mensaje anade overhead de CPU y puede dedup incorrectamente para casos mismo-payload-diferente-intencion.

Solución

Python (Redis + SQS)

import redis
import json
import hashlib

r = redis.Redis(host="localhost", port=6379, db=0)
DEDUP_TTL = 86400  # 24 horas

def process_message(message_id, payload):
    # Verificar si ya fue procesado
    dedup_key = f"dedup:{message_id}"
    if r.exists(dedup_key):
        print(f"Skipping duplicate message {message_id}")
        return

    # Procesar el mensaje
    result = handle_order(payload)

    # Marcar como procesado con TTL
    r.setex(dedup_key, DEDUP_TTL, "1")
    print(f"Processed message {message_id}")

def handle_order(payload):
    order = json.loads(payload)
    print(f"Charging payment for order {order['order_id']}")
    return {"status": "charged"}

# Simular entrega duplicada
process_message("msg-001", '{"order_id": 42}')
process_message("msg-001", '{"order_id": 42}')  # Saltado

JavaScript (Redis + BullMQ)

import Redis from "ioredis";
import { Worker } from "bullmq";

const redis = new Redis({ host: "localhost", port: 6379 });
const DEDUP_TTL = 86400; // 24 horas

async function isDuplicate(messageId) {
  const key = `dedup:${messageId}`;
  const result = await redis.set(key, "1", "EX", DEDUP_TTL, "NX");
  // result es "OK" si la clave fue establecida (primera vez), null si ya existia
  return result === null;
}

const worker = new Worker(
  "orders",
  async (job) => {
    const messageId = job.data.messageId;
    const payload = job.data.payload;

    if (await isDuplicate(messageId)) {
      console.log(`Skipping duplicate message ${messageId}`);
      return { status: "skipped" };
    }

    // Procesar el mensaje
    console.log(`Charging payment for order ${payload.orderId}`);
    return { status: "processed" };
  },
  { connection: { host: "localhost", port: 6379 } }
);

Java (Redis + Spring)

import org.springframework.data.redis.core.StringRedisTemplate;
import org.springframework.stereotype.Component;
import java.util.concurrent.TimeUnit;

@Component
public class DeduplicatingConsumer {

    private final StringRedisTemplate redis;
    private static final int DEDUP_TTL = 86400; // 24 horas

    public DeduplicatingConsumer(StringRedisTemplate redis) {
        this.redis = redis;
    }

    public void processMessage(String messageId, String payload) {
        String dedupKey = "dedup:" + messageId;

        // Check-and-set atomico: retorna true si la clave fue establecida (primera vez)
        Boolean wasSet = redis.opsForValue()
            .setIfAbsent(dedupKey, "1", DEDUP_TTL, TimeUnit.SECONDS);

        if (Boolean.FALSE.equals(wasSet)) {
            System.out.println("Skipping duplicate message " + messageId);
            return;
        }

        // Procesar el mensaje
        System.out.println("Processing message " + messageId + ": " + payload);
        handleOrder(payload);
    }

    private void handleOrder(String payload) {
        // Logica de negocio aqui
    }
}

Explicación

El patron usa una operacion atomica check-and-set (Redis SET NX EX) para determinar si un mensaje fue procesado. La operacion es atomica: solo un consumidor puede establecer la clave, por lo que consumidores concurrentes procesando el mismo mensaje no procederan ambos.

El TTL en la clave de deduplicacion evita que el almacen crezca indefinidamente. Tras expirar el TTL, el mismo ID de mensaje podria procesarse de nuevo. Elige un TTL mayor que tu ventana maxima de reentrega (tipicamente 24 horas para SQS, o el periodo de retencion de la cola).

Para que la deduplicacion funcione, cada mensaje debe llevar un identificador unico. Puede ser un hash de contenido (para deduplicacion basada en payload) o un ID asignado por el productor (para deduplicacion basada en identidad).

Variantes

VarianteAlmacenCaso de UsoCompromiso
Redis SET NXRedisRapido, compartido entre consumidoresDependencia externa, expiracion por TTL
Restriccion Unica DBSQL DBDurable, transaccionalMas lento, anade carga a DB
Hash de ContenidoCualquier almacenDedup por contenido de payloadMismo payload siempre se dedup, incluso si diferente intencion
Nativo del BrokerSQS FIFODedup integradoSolo colas FIFO, throughput limitado
Set en MemoriaMemoria del procesoConsumidor unico, rapidoSe pierde al reiniciar, no compartido entre instancias

Qué Funciona

  • Usa check-and-set atomico (Redis SET NX EX) para evitar race conditions entre consumidores concurrentes
  • Establece TTL mayor que la ventana de reentrega del broker
  • Usa hashes de contenido cuando los mensajes carecen de IDs explicitos
  • Haz consumidores idempotentes como segunda capa de defensa: incluso si dedup falla, procesar dos veces no debe causar dano
  • Registra duplicados saltados para debugging y monitoreo
  • Usa colas FIFO con deduplicacion basada en contenido cuando el broker lo soporta (SQS FIFO)

Errores Comunes

  • Check-then-set sin atomicidad: Dos consumidores verifican simultaneamente, ambos ven que no hay clave, ambos procesan.
  • TTL demasiado corto: Si el TTL expira antes de que el broker deje de reentregar, los duplicados pasan. Establece TTL al menos 2x la ventana de reentrega.
  • Usar memoria del proceso para dedup: Se pierde al reiniciar, no se comparte entre instancias.
  • No manejar fallos del almacen de dedup: Si Redis cae, dedup falla. Decide si procesar (riesgo de duplicados) o rechazar (perder mensajes).
  • Clave de dedup basada en campos mutables: Si la clave incluye campos que cambian, el mismo mensaje logico obtiene claves diferentes y se procesa dos veces.

Como Funciona

  1. El mensaje llega con ID unico: Cada mensaje lleva una clave de deduplicacion — ya sea un UUID asignado por el productor o un hash de contenido. El consumidor extrae esta clave antes de procesar.
  2. Check-and-set atomico: El consumidor intenta SET NX EX en la clave de dedup en un almacen compartido (Redis). Si la clave no existe, el consumidor la establece y procede. Si existe, el mensaje ya fue procesado — saltalo.
  3. Procesa el mensaje: Solo el consumidor que establecio la clave exitosamente procesa el mensaje. Consumidores concurrentes con el mismo ID de mensaje fallan el SET NX y saltan.
  4. Expiracion TTL: Tras el TTL configurado, la clave expira. Esto previene que el almacen crezca indefinidamente. El TTL debe exceder la ventana maxima de reentrega del broker.

La atomicidad de SET NX es critica: sin ella, dos consumidores podrian verificar, ambos ver que no hay clave, y ambos procesar el mismo mensaje.

Mejores Practicas

  • Usa IDs asignados por el productor sobre hashes de contenido. Los hashes dedup payloads identicos, lo cual puede ser incorrecto si el mismo payload representa operaciones logicas diferentes. Los IDs del productor dedup por identidad, no contenido.
  • Loggea duplicados saltados. Cuando un consumidor salta un duplicado, loggea el ID del mensaje y timestamp. Ayuda a debuggear problemas de reentrega y medir tasas de duplicados.
  • Monitorea la tasa de hits de dedup. Si 50% de mensajes son duplicados, algo esta mal upstream — el productor esta reintentando demasiado agresivamente o el consumidor es muy lento para acusar recibo.
  • Usa colas FIFO cuando esten disponibles. SQS FIFO y partition keys de Kafka proporcionan ordenamiento y dedup a nivel del broker, reduciendo la necesidad de dedup a nivel aplicacion.
  • Degradacion elegante ante fallo del almacen de dedup. Si Redis cae, decide entre procesar con idempotencia o rechazar. Documenta la eleccion y alerta sobre ello.

Ejemplos del Mundo Real

Webhooks de Pago de Stripe

Stripe envia webhooks para eventos de pago. Fallos de red causan que Stripe reintente, entregando el mismo webhook multiples veces. Stripe recomienda usar IDs de evento para deduplicacion: verifica el ID de evento en un almacen antes de procesar. Sin dedup, un solo pago podria disparar multiples fulfillments de pedido.

SQS + Lambda para Procesamiento de Pedidos

Una plataforma de e-commerce usa SQS para disparar Lambda para procesamiento de pedidos. SQS puede reentregar mensajes si Lambda no acusa recibo a tiempo. La plataforma usa Redis SET NX con el ID del pedido como clave de dedup. Entregas duplicadas se saltan, previniendo doble cobro o doble envio.

Consumidor Kafka con Dedup en Redis

Un pipeline de streaming consume eventos de Kafka y escribe a una base de datos. La entrega at-least-once de Kafka significa que el mismo evento puede consumirse dos veces. El consumidor verifica Redis antes de escribir a la base de datos. Esto proporciona semantica exactly-once sin transacciones de Kafka.

Referencia Rápida

  • Comando principal: ejecuta la solución base del artículo y verifica el resultado esperado.
  • Validación: confirma que los tests pasan y que las métricas clave no se degradaron.
  • Rollback: si algo falla, revierte el cambio y consulta la sección de Troubleshooting.

Lectura Adicional

  • Documentación oficial: consulta la referencia actualizada del framework o herramienta utilizada.
  • Guías relacionadas: explora las guías de deduplication y pattern para profundizar.
  • Patrones complementarios: revisa los patrones de diseño aplicables a tu stack tecnológico.
  • Postmortems públicos: estudia incidentes reales de equipos que enfrentaron problemas similares en producción.

Notas de Producción

  • Despliega gradualmente usando canary o blue-green para detectar regresiones temprano.
  • Configura alertas para errores, latencia p99 y tasa de fallos antes de habilitar en producción.
  • Documenta el rollback en el runbook; prueba el procedimiento en staging al menos una vez por trimestre.
  • Revisa logs estructurados con correlation IDs para trazar requests end-to-end en incidentes.

Puntos Clave

  • Aplica patrón message deduplication cuando necesites una solución práctica para tu caso de uso.
  • Monitorea el rendimiento después de implementar; mide latencia, errores y uso de recursos antes y después.
  • Revisa la sección de Troubleshooting ante errores comunes; la mayoría tienen causa raíz documentada con solución.
  • Mantén dependencias actualizadas y ejecuta tests en CI para prevenir regresiones en producción.

Troubleshooting

  • Messages are lost on restart: persist messages before acknowledging.
  • Consumer lags behind producer: scale consumers, increase prefetch, and partition the topic.
  • Duplicate messages: design consumers to be idempotent.
  • Ordering is wrong after scaling: preserve partition keys and avoid rebalancing during bursts. Consider a single partition when order is mandatory.
  • Queue depth grows but consumers are idle: check network partitions, consumer health, and permission issues. Restart gracefully.

Errores Comunes en Producción

  • Aplicar el patrón donde no se necesita abstracción, agregando complejidad accidental.
  • Dejar que el patrón se filtre en módulos no relacionados y confundir los límites de responsabilidad.
  • Sobre-ingeniería en la primera implementación en lugar de comenzar simple y medir el dolor.
  • Saltar los tests de contrato, de modo que las refactorizaciones rompan consumidores en silencio.
  • Ignorar modos de fallo que el patrón no cubre.
  • Usar el patrón como opción por defecto en lugar de elegir la herramienta adecuada para la escala actual.
  • Olvidar documentar cuándo dejar de usar el patrón y qué lo reemplaza.
  • Carecer de observabilidad sobre rendimiento y propagación de errores del patrón.