Artículo

Laravel y Kafka sin eventos perdidos: outbox transaccional y consumidores idempotentes con PostgreSQL

Laravel y Kafka sin eventos perdidos: outbox transaccional y consumidores idempotentes con PostgreSQL

Una confirmación en PostgreSQL y una publicación en Kafka no forman una transacción ordinaria. Aprende a evitar pérdidas y duplicación de efectos.

11 min de lecturaIdioma: ES EspañolGratis0 aplausos0 comentarios
Opciones de lectura

Esta implementación parece sencilla:

$order = Order::create($data);
$kafka->publish('orders.placed.v1', $order->toArray());

Si la inserción funciona y la publicación falla, el pedido existe pero inventario nunca lo recibe. Invertir las llamadas puede publicar un evento de un pedido que después se revierte. Meter la llamada de red dentro de DB::transaction() tampoco incorpora Kafka a la transacción de PostgreSQL.

El outbox transaccional guarda el pedido y un registro de evento en la misma transacción de base de datos. Otro proceso publica el registro. Así no se pierde el evento si el proceso web muere tras confirmar. La publicación puede repetirse, por lo que el consumidor también debe evitar efectos duplicados.

1. Define el contrato del evento

{
  "eventId": "evt_01JQ8X8F",
  "eventType": "OrderPlaced",
  "schemaVersion": 1,
  "occurredAt": "2026-09-29T10:30:00Z",
  "orderId": "ord_782",
  "customerId": "cus_17",
  "totalMinor": 12900,
  "currency": "USD"
}

OrderPlaced significa que el pedido ya está confirmado. eventId permanece igual en todos los reintentos y sirve para deduplicar. La clave Kafka es orderId, útil para conservar el orden de eventos de ese pedido bajo una partición estable. No publiques el modelo Eloquent completo: columnas, relaciones y datos internos no son un contrato público.

Documenta propietario del tema, clave, campos, retención, compatibilidad y consumidores previstos. El sufijo v1 no sustituye una política de esquemas.

2. Una tabla outbox y una transacción

Un punto de partida en PostgreSQL:

CREATE TABLE outbox_events (
  id uuid PRIMARY KEY,
  aggregate_id text NOT NULL,
  topic text NOT NULL,
  event_key text NOT NULL,
  payload jsonb NOT NULL,
  occurred_at timestamptz NOT NULL,
  published_at timestamptz NULL,
  attempts integer NOT NULL DEFAULT 0
);
CREATE INDEX outbox_pending_idx ON outbox_events (occurred_at, id)
  WHERE published_at IS NULL;

En el mismo DB::transaction() crea pedido y fila outbox. Si hay rollback, desaparecen ambos. Si hay commit, el evento queda disponible aunque el proceso web se detenga. Mantén corta la transacción; no esperes la red de Kafka mientras retienes bloqueos. Esquema de código adaptable:

DB::transaction(function () use ($data) {
    $order = Order::create($data);
    $eventId = (string) Str::uuid();
    DB::table('outbox_events')->insert([
        'id' => $eventId,
        'aggregate_id' => (string) $order->id,
        'topic' => 'orders.placed.v1',
        'event_key' => (string) $order->id,
        'payload' => json_encode([
            'eventId' => $eventId,
            'eventType' => 'OrderPlaced',
            'schemaVersion' => 1,
            'orderId' => (string) $order->id,
        ], JSON_THROW_ON_ERROR),
        'occurred_at' => now(),
    ]);
});

Valida entrada y tipos, minimiza datos personales y modela el contrato sin depender de la estructura de la tabla de pedidos. El ejemplo no pretende ser una API universal de un paquete Kafka.

3. El publicador y su ventana de fallo

Un worker selecciona pocas filas pendientes, las reclama de forma segura entre varios workers, publica usando event_key, espera la confirmación del productor y marca published_at. Usa una reclamación breve o un lease recuperable; evita mantener un bloqueo de fila mientras un broker responde lentamente.

Kafka puede aceptar el evento y el worker morir antes de marcar la fila. Otro worker lo publicará de nuevo. El outbox ofrece publicación al menos una vez, no un efecto global exactamente una vez. Conserva el mismo eventId al reintentar. Un productor idempotente reduce duplicados derivados de sus reintentos hacia Kafka, pero no vuelve atómica la actualización PostgreSQL con la confirmación Kafka.

Puedes implementarlo con un comando Laravel, worker o CDC. Evalúa cliente, compatibilidad y seguridad. Los drivers integrados de Laravel incluyen Redis, SQS y base de datos; Kafka no aparece mágicamente como driver estándar al cambiar QUEUE_CONNECTION.

4. Evita reservar dos veces en el consumidor

Inventario puede reservar stock y caer antes de confirmar su offset. Entonces leerá otra vez el evento. Añade una restricción única:

CREATE TABLE processed_events (
  consumer_name text NOT NULL,
  event_id text NOT NULL,
  processed_at timestamptz NOT NULL DEFAULT now(),
  PRIMARY KEY (consumer_name, event_id)
);

Dentro de una transacción de la base de inventario, inserta (inventory-reservation, eventId) y aplica la reserva. Si ya existe, omite el efecto. Solo después del commit avanza el offset Kafka. Si la base confirma y el offset no, una reproducción será segura. Si la base revierte, el offset no debe avanzar como si hubiera éxito.

El correo o un proveedor de pagos no participa en esa transacción local. Usa la clave de idempotencia del proveedor si existe, otro outbox o un plan compensatorio. Explica la garantía de cada efecto, no una supuesta garantía universal del pipeline.

5. Reintentos, cuarentena y evolución

Una caída temporal de base de datos admite reintento con espera. Un evento malformado persistente requiere inspección, corrección o cuarentena. Un tema dead-letter es una opción, pero publicar allí y confirmar el offset original tiene otra ventana de fallo; documenta secuencia, alertas, acceso, retención y replay. No captures toda excepción para confirmar el offset sin completar el trabajo.

Registra eventId, tema, partición, offset, grupo y categoría segura de error; evita payloads personales completos. Conserva deduplicación durante la ventana prometida de retención y reproducción. Prueba consumidores contra eventos antiguos al evolucionar el esquema y no cambies silenciosamente el significado de un campo.

Pruebas que importan

Prueba rollback del pedido, commit con publicador detenido, publicación aceptada y caída antes de published_at, doble entrega a inventario, caída después de la reserva y antes del offset, y workers concurrentes. Vigila edad de la fila pendiente más antigua, pendientes, errores de publicación, lag y cuarentena. Un HTTP saludable no garantiza que el outbox esté drenando.

Una cola Laravel normal basta para muchos trabajos sencillos. Si varios sistemas independientes dependen del hecho de negocio y el replay aporta valor, el diseño fiable reúne transacción, evento durable, identidad estable, consumidor idempotente y fallos observables.

Fuentes

Artículos destacados

Kafka en producción: particiones, lag, fiabilidad y respuesta a incidentes
EditorialES
12 minGratis

Kafka en producción: particiones, lag, fiabilidad y respuesta a incidentes

Una guía práctica para tomar decisiones de orden, capacidad y recuperación en Kafka, interpretar el lag y responder a fallos sin perder de vista el resultado de negocio.

Artículos de ingenieríaGuías de plataforma
0 aplausos
Leer
Apache Kafka explicado: temas, particiones, grupos de consumidores y tu primer flujo de eventos
EditorialES
10 minGratis

Apache Kafka explicado: temas, particiones, grupos de consumidores y tu primer flujo de eventos

Sigue un evento de pedido desde el productor hasta varios consumidores y prueba en local cómo funcionan las particiones, el orden y la reproducción.

Artículos de ingenieríaGuías de plataforma
0 aplausos
Leer
Cómo diseñar APIs en las que los clientes puedan confiar: guía práctica desde el contrato
EditorialES
10 minGratis

Cómo diseñar APIs en las que los clientes puedan confiar: guía práctica desde el contrato

Una buena API hace predecible la próxima acción del cliente. Diseña una inscripción a un curso con reintentos seguros, errores útiles y un plan de evolución.

Artículos de ingenieríaGuías de plataforma
0 aplausos
Leer

Comentarios

0 comentarios

Todavía no hay comentarios aprobados. Las respuestas nuevas pueden esperar moderación.