intermediate Por Mathias Paulenko

Patron de Convoy Secuencial

Preserva el orden de mensajes relacionados en un sistema distribuido agrupandolos en secuencias ordenadas y procesandolos uno a la vez a traves de un unico consumidor.

Nota para desarrolladores hispanohablantes: Esta guía incluye ejemplos y convenciones de nomenclatura adaptadas a equipos que trabajan en español. Cuando existen diferencias significativas en terminología técnica entre el inglés y el español, se indican explícitamente para facilitar la comunicación en equipos multiculturales.

Resumen

El Patron de Convoy Secuencial preserva el orden de mensajes relacionados en un sistema de mensajeria distribuido. Cuando los mensajes tienen una relacion causal, procesarlos fuera de orden produce estado inconsistente.

Este patron agrupa mensajes relacionados en un convoy (secuencia) y asegura que sean procesados por un unico consumidor en orden. Convoys no relacionados pueden procesarse en paralelo, preservando correccion y rendimiento.

Cuando Usar

  • For alternatives, see Idempotent Consumer Pattern.

  • Mensajes para la misma entidad deben procesarse en orden de produccion

  • Event sourcing donde eventos para un agregado deben aplicarse secuencialmente

  • Pipelines de procesamiento de pedidos donde transiciones dependen de estados previos

  • Sistemas de inventario donde movimientos de stock deben aplicarse cronologicamente

  • Workflows multi-paso donde el paso N no puede comenzar hasta que el paso N-1 termine

Cuando Evitar

  • Mensajes sin relacion causal — el procesamiento paralelo es mas simple y rapido
  • El orden estricto no es necesario (ej. eventos de analytics independientes)
  • El sistema tolera consistencia eventual sin garantias de orden
  • Volumenes por convoy son tan altos que un solo consumidor crea cuello de botella

Solucion

Python (Kafka con Partition Key)

from kafka import KafkaProducer, KafkaConsumer
import json
import time

class OrderedMessageProducer:
    """Produce mensajes que mantienen ordenamiento por entidad"""

    def __init__(self, bootstrap_servers):
        self.producer = KafkaProducer(
            bootstrap_servers=bootstrap_servers,
            value_serializer=lambda v: json.dumps(v).encode('utf-8'),
            partitioner=lambda key, partitions, topic: (
                hash(key) % len(partitions) if key else 0
            )
        )

    def send_event(self, entity_id: str, event_type: str, payload: dict):
        message = {
            'entity_id': entity_id,
            'event_type': event_type,
            'payload': payload,
            'timestamp': time.time(),
            'sequence_number': self._get_next_sequence(entity_id)
        }
        self.producer.send(
            'entity-events',
            key=entity_id.encode('utf-8'),
            value=message
        )
        self.producer.flush()

    def _get_next_sequence(self, entity_id: str) -> int:
        return int(time.time() * 1000)

class SequentialConvoyConsumer:
    """Procesa mensajes en orden por entidad"""

    def __init__(self, bootstrap_servers, topic, group_id):
        self.consumer = KafkaConsumer(
            topic, bootstrap_servers=bootstrap_servers,
            group_id=group_id, auto_offset_reset='earliest',
            max_poll_records=1, enable_auto_commit=False
        )
        self.pending_sequences: dict = {}
        self.last_processed: dict = {}

    def process_messages(self):
        for message in self.consumer:
            event = json.loads(message.value)
            entity_id = event['entity_id']
            seq_num = event['sequence_number']
            expected_seq = self.last_processed.get(entity_id, 0) + 1

            if seq_num == expected_seq:
                self._process_event(event)
                self.last_processed[entity_id] = seq_num
                self._check_pending(entity_id)
            elif seq_num > expected_seq:
                self.pending_sequences.setdefault(entity_id, {})[seq_num] = event
                print(f"Buffering mensaje fuera de orden {seq_num} para {entity_id}")
            else:
                print(f"Saltando duplicado {seq_num} para {entity_id}")

            self.consumer.commit()

    def _process_event(self, event):
        print(f"Procesando {event['event_type']} para {event['entity_id']}")

    def _check_pending(self, entity_id):
        pending = self.pending_sequences.get(entity_id, {})
        expected = self.last_processed.get(entity_id, 0) + 1
        while expected in pending:
            event = pending.pop(expected)
            self._process_event(event)
            self.last_processed[entity_id] = expected
            expected += 1

# Uso
producer = OrderedMessageProducer(['localhost:9092'])
producer.send_event('user-123', 'created', {'name': 'Alice'})
producer.send_event('user-123', 'updated', {'name': 'Alice Smith'})
producer.send_event('user-123', 'deleted', {})

consumer = SequentialConvoyConsumer(['localhost:9092'], 'entity-events', 'convoy-group')
consumer.process_messages()

Java (Azure Service Bus Sessions)

import com.azure.messaging.servicebus.ServiceBusClientBuilder;
import com.azure.messaging.servicebus.ServiceBusMessage;
import com.azure.messaging.servicebus.ServiceBusSenderClient;
import com.azure.messaging.servicebus.ServiceBusProcessorClient;
import com.azure.messaging.servicebus.models.ServiceBusReceiveMode;
import java.util.UUID;
import java.util.concurrent.ConcurrentHashMap;

public class SequentialConvoyServiceBus {

    private static final String CONNECTION_STRING = "<connection-string>";
    private static final String QUEUE_NAME = "ordered-queue";

    public static class OrderedProducer {
        private final ServiceBusSenderClient sender;

        public OrderedProducer() {
            this.sender = new ServiceBusClientBuilder()
                .connectionString(CONNECTION_STRING)
                .sender()
                .queueName(QUEUE_NAME)
                .buildClient();
        }

        public void sendOrderedEvents(String entityId, List<DomainEvent> events) {
            for (int i = 0; i < events.size(); i++) {
                ServiceBusMessage message = new ServiceBusMessage(
                    serializeEvent(events.get(i))
                );
                message.setSessionId(entityId);
                message.setApplicationProperty("sequenceNumber", i);
                sender.sendMessage(message);
            }
        }
    }

    public static class OrderedConsumer {

        public void startProcessing() {
            ServiceBusProcessorClient processor = new ServiceBusClientBuilder()
                .connectionString(CONNECTION_STRING)
                .processor()
                .queueName(QUEUE_NAME)
                .receiveMode(ServiceBusReceiveMode.PEEK_LOCK)
                .processMessage(this::processMessage)
                .processError(this::handleError)
                .prefetchCount(1)
                .buildProcessorClient();

            processor.start();
        }

        private void processMessage(ServiceBusReceivedMessageContext context) {
            ServiceBusReceivedMessage message = context.getMessage();
            String sessionId = message.getSessionId();
            int sequenceNumber = (int) message.getApplicationProperties()
                .get("sequenceNumber");

            DomainEvent event = deserializeEvent(message.getBody().toString());
            applyEvent(sessionId, event);
            context.complete();
        }

        private void handleError(ServiceBusErrorContext context) {
            System.err.println("Error: " + context.getException().getMessage());
        }
    }
}

JavaScript (Redis Streams con Consumer Groups)

const Redis = require('ioredis');

class SequentialConvoyProcessor {
    constructor(redis, streamKey) {
        this.redis = redis;
        this.streamKey = streamKey;
        this.groupName = 'convoy-processors';
    }

    async initialize() {
        try {
            await this.redis.xgroup('CREATE', this.streamKey,
                this.groupName, '0', 'MKSTREAM');
        } catch (err) {
            if (!err.message.includes('already exists')) throw err;
        }
    }

    async produceEvent(entityId, eventType, payload) {
        const sequence = await this.redis.incr(`seq:${entityId}`);
        await this.redis.xadd(this.streamKey, '*',
            'entityId', entityId,
            'sequence', sequence.toString(),
            'eventType', eventType,
            'payload', JSON.stringify(payload)
        );
    }

    async consumeOrdered(consumerId) {
        const results = await this.redis.xreadgroup(
            'GROUP', this.groupName, consumerId,
            'COUNT', 1,
            'BLOCK', 5000,
            'STREAMS', this.streamKey, '>'
        );

        if (!results || results.length === 0) return null;

        const [[, messages]] = results;
        const [id, fields] = messages[0];
        const event = this.parseFields(fields);
        const entityId = event.entityId;
        const sequence = parseInt(event.sequence);

        const lastProcessed = await this.redis.get(`last:${entityId}`);
        const expected = lastProcessed ? parseInt(lastProcessed) + 1 : 1;

        if (sequence === expected) {
            await this.processEvent(event);
            await this.redis.set(`last:${entityId}`, sequence);
            await this.redis.xack(this.streamKey, this.groupName, id);
            return event;
        } else if (sequence > expected) {
            console.log(`Fuera de orden: esperado ${expected}, recibido ${sequence}`);
            return null;
        } else {
            await this.redis.xack(this.streamKey, this.groupName, id);
            return null;
        }
    }

    parseFields(fields) {
        const obj = {};
        for (let i = 0; i < fields.length; i += 2) {
            obj[fields[i]] = fields[i + 1];
        }
        return obj;
    }

    async processEvent(event) {
        console.log(`Procesando ${event.eventType} para ${event.entityId}`);
    }
}

Explicacion

El patron se basa en dos mecanismos clave:

  1. Particionamiento por ID de entidad: Los mensajes para la misma entidad se enrutan a la misma particion/cola/sesion. Esto se hace usando partition key (Kafka), session ID (Service Bus), o un campo entity (Redis).

  2. Unico consumidor por particion: Solo un consumidor procesa mensajes de una particion dada a la vez. Esto previene que dos consumidores manejen diferentes mensajes para la misma entidad simultaneamente, lo que violaria el orden.

El compromiso es paralelismo reducido por entidad — todos los mensajes para user-123 deben procesarse secuencialmente. Sin embargo, mensajes para user-456 pueden procesarse en paralelo en otra particion.

Variantes

VarianteMecanismoIdeal Para
Partition key de KafkaAsignacion por hashAlto throughput, orden simple
Sesiones de Service BusBalanceo por sesionNativo en la nube, exactly-once por sesion
Single active consumer de RabbitMQConsumidor exclusivo por colaOrdenamiento basado en colas simple
Tabla de secuencia en base de datosBloqueo optimista en sequence numbersSistemas sin ordering del broker
Sagas con orquestacionOrden explicito de pasos en workflow engineProcesos de negocio multi-paso complejos

Lo que Funciona

  • Usar una clave de particion determinista. El ID de entidad debe mapear consistentemente a la misma particion. Cambiar la clave invalida el orden.
  • Monitorear skew de particiones. Si una entidad genera 90% de los mensajes, su particion se vuelve un cuello de botella. Considera dividir entidades hot.
  • Manejar mensajes faltantes con gracia. Si la secuencia N nunca llega, el convoy se detiene. Implementa timeouts y alertas.
  • Mantener convoys pequenos. Convoys largos retienen mensajes nuevos. Disena para secuencias cortas y acotadas.
  • Procesamiento idempotente dentro de convoys. Incluso con orden, los reintentos pueden causar duplicados. Haz las operaciones individuales idempotentes.

Errores Comunes

  • Cambiar claves de particion. Rebalancear particiones de Kafka cambia que consumidor maneja que entidad, violando supuestos de orden.
  • Multiples consumidores por particion. Dos consumidores leyendo la misma particion procesaran mensajes en paralelo para la misma entidad.
  • No manejar gaps de secuencia. Un mensaje perdido en una secuencia bloquea todos los mensajes subsiguientes para siempre.
  • Convoys demasiado grandes. Un convoy que procesa miles de mensajes para una entidad crea un hotspot.
  • Ignorar reintentos del productor. Un mensaje reintentado puede reordenarse relativo a un mensaje mas nuevo si no usa la misma partition key.

Ejemplos del Mundo Real

Particionamiento de Kafka

Kafka garantiza orden dentro de una particion. Usando el ID de usuario como partition key, todos los eventos para un usuario estan ordenados. Uber usa esto para eventos de viaje: trip-created, driver-assigned, trip-started, trip-completed deben procesarse en orden para calculo de tarifa.

Sesiones de Azure Service Bus

Las sesiones de Service Bus proporcionan orden FIFO dentro de una sesion. Una plataforma de e-commerce usa sesiones por carrito de compras: item-added, quantity-changed, checkout-initiated, payment-received deben procesarse secuencialmente para mantener consistencia del carrito.

Event Store DB

Event Store DB usa control de concurrencia optimista en streams. Cada agregado (ej. un pedido) es un stream, y los eventos se agregan con versiones esperadas. Escritores concurrentes fallan si el stream fue modificado, preservando orden.