Zum Inhalt springen
Architektur

Transactional inbox: cuando el ack borra el mensaje

Un ack a Google Cloud Pub/Sub no es un comprobante de procesamiento, sino una orden de borrado. Si entre la recepción y el ack hay un motor de workflow que también puede rechazar un mensaje por motivos de negocio, el lugar de esa única llamada decide si una notificación que no llega es un incidente operativo visible o una pérdida de datos silenciosa. Este artículo es un tutorial completo: construye el receptor típico, desmonta sus modos de fallo y le contrapone cinco variantes, desde nack con Dead-Letter-Topic hasta la Transactional Inbox con relay. Todo el código corre sobre Spring Boot con Spring Cloud GCP y Flowable.

Contenido

La escena: una notificación que nunca llega

Un sistema de órdenes de trabajo entrega órdenes a un proveedor de servicios que las ejecuta. Una orden comprende varios miles de posiciones. El flujo lo dirige un proceso BPMN work-order-execution: entrega la orden al proveedor, espera después en un Intermediate Message Catch Event su aviso de finalización «orden de trabajo completada», y solo entonces la orden se da por terminada y se factura.

┌───────────────────┐   ┌─────────────────────────┐   ┌─────────────────────┐
│ Entregar la orden │───│ Esperar la finalización │───│ Emitir la factura   │
└───────────────────┘   └─────────────────────────┘   └─────────────────────┘
      Send Task             Message Catch Event         Service Task, async

El aviso de finalización llega por Google Cloud Pub/Sub. El proveedor publica en un topic, un servicio Spring Boot llamado work-order-service mantiene la subscription y correlaciona cada mensaje con la instancia de proceso en curso a través del Correlation Key workOrderId.

Si el aviso de finalización no llega, no hay error, ni alarma roja, ni stacktrace. El proceso se queda parado en el Catch Event esperando, la orden sigue abierta y no se factura nada. Por experiencia, un parón así solo se detecta cuando alguien pregunta por el paradero de una orden, y entonces no se trata de una posición, sino de varios miles.

La pregunta de este tutorial es por tanto: ¿cómo construyes el camino del mensaje de Pub/Sub a la instancia de proceso de forma que una avería en ese camino siga siendo visible y ningún mensaje se pierda de forma definitiva?

El receptor que parece inofensivo

Así es el receptor que probablemente todo el mundo ha escrito alguna vez. Usa PubSubTemplate.subscribe de Spring Cloud GCP (starter com.google.cloud:spring-cloud-gcp-starter-pubsub, versión 8.1.0 para Spring Boot 4.0 y 4.1, junto a la 7.4.10 para Spring Boot 3.5), lee el mensaje, lo correlaciona con el motor y hace ack.

@Component
public class WorkOrderCompletedReceiver {

    private static final Logger log =
            LoggerFactory.getLogger(WorkOrderCompletedReceiver.class);

    private final PubSubTemplate pubSubTemplate;
    private final WorkOrderProcessService processService;
    private final ObjectMapper objectMapper;

    public WorkOrderCompletedReceiver(PubSubTemplate pubSubTemplate,
                                      WorkOrderProcessService processService,
                                      ObjectMapper objectMapper) {
        this.pubSubTemplate = pubSubTemplate;
        this.processService = processService;
        this.objectMapper = objectMapper;
    }

    @PostConstruct
    void startSubscription() {
        pubSubTemplate.subscribe("work-order-completed-sub", this::process);
    }

    private void process(BasicAcknowledgeablePubsubMessage message) {
        try {
            WorkOrderCompleted event = objectMapper.readValue(
                    message.getPubsubMessage().getData().toByteArray(),
                    WorkOrderCompleted.class);

            processService.workOrderCompleted(event.workOrderId(), Map.of(
                    "completedItems", event.completedItems(),
                    "completionTimestamp", event.completionTimestamp().toString()));

            message.ack();
        } catch (Exception e) {
            log.error("Notificación no procesable", e);
            message.ack();   // para que la subscription no se atasque
        }
    }
}
public record WorkOrderCompleted(String workOrderId,
                                 int completedItems,
                                 Instant completionTimestamp) {
}

El código compila, corre, pasa cualquier test de demo y sobrevive meses en producción. Aun así tiene tres debilidades silenciosas, y cada una recibe su propio capítulo en este artículo:

  1. El motor de workflow está en la ruta del ack. Cualquier avería del motor, incluso un deployment normal, impide el ack y convierte al receptor en un cuello de botella.
  2. Todos los errores caen en el mismo catch. Una conexión de base de datos rota y un rechazo de negocio de la correlación se tratan igual, aunque necesitarían reacciones opuestas.
  3. El ack está en el bloque catch. El comentario de al lado suena razonable; en realidad esa línea es una orden de borrado para todo mensaje que haya provocado un error.

Las tres dependen de lo que un ack significa realmente en Pub/Sub.

Qué promete realmente un ack

Un ack a Pub/Sub promete exactamente una cosa: este mensaje no necesita más entregas. En cuanto, para cada subscription, al menos un subscriber ha confirmado el mensaje, Pub/Sub lo borra de su almacenamiento. Si tu código procesó el mensaje con éxito antes, Pub/Sub no lo sabe y no lo comprueba.

Tres relojes determinan el comportamiento:

La ack deadline decide sobre la reentrega. Si no confirmas un mensaje dentro de la deadline, Pub/Sub lo entrega de nuevo. El default está en 10 segundos, configurable entre 10 y 600 segundos. La client library va extendiendo la deadline en segundo plano mientras el mensaje sigue en tus manos, con el límite de la maxAckExtensionPeriod del cliente.

La message retention de la subscription decide la fecha de caducidad definitiva. El default son 7 días, configurable entre 10 minutos y 31 días. Pasado ese plazo, Pub/Sub puede descartar el mensaje, con independencia de si fue confirmado o no. Quien no pudo procesar un mensaje durante 31 días lo pierde también sin ningún ack.

Y el tercer reloj ni siquiera existe: un mensaje confirmado no tiene vuelta atrás. Sin la opción «Retain acknowledged messages» en la subscription o una message retention configurada en el topic, un mensaje con ack no se puede recuperar por seek. Por eso el ack del bloque catch del receptor de arriba es definitivo.

La garantía de entrega que acompaña a esto: at-least-once es el default para todos los tipos de subscription. Cualquier mensaje puede llegar varias veces, incluso sin que nadie haya cometido un error, por ejemplo porque un ack se perdió por el camino en la red. Un receptor que no tolera duplicados está mal construido para Pub/Sub. Qué pasa con exactly-once delivery lo aclara la FAQ; un adelanto: no resuelve el problema de este artículo.

Por qué la correlación rechaza

Del lado del motor, el proceso espera en un Message Catch Event. En BPMN se ve así:

<message id="workOrderCompletedMessage" name="WorkOrderCompleted" />

<process id="work-order-execution" name="Ejecución de la orden de trabajo">

    <sendTask id="handOverWorkOrder" name="Entregar la orden de trabajo"
              flowable:async="true" />
    <sequenceFlow sourceRef="handOverWorkOrder" targetRef="waitForCompletion" />

    <intermediateCatchEvent id="waitForCompletion" name="Esperar la finalización">
        <messageEventDefinition messageRef="workOrderCompletedMessage" />
    </intermediateCatchEvent>
    <sequenceFlow sourceRef="waitForCompletion" targetRef="createInvoice" />

    <serviceTask id="createInvoice" name="Emitir la factura"
                 flowable:async="true" />
</process>

Tres elementos, y el del medio es el vulnerable. El Send Task entrega la orden, el Catch Event espera, el Service Task emite la factura. Mientras ningún token esté en el Catch Event, no existe receptor para el aviso de finalización.

Flowable no ofrece en el API core una llamada de un solo paso que entregue un mensaje por sí misma a partir de un Correlation Key. El llamador construye la query, evalúa los resultados, entrega y maneja los errores. Estas consultas se llaman correlation queries, y quien escribe una asume con ella una hipótesis sobre el modelo de proceso: aquí, que por cada workOrderId espera como máximo una instancia. La lógica de correlación vive en tu receiver, no en el motor.

@Service
public class WorkOrderProcessService {

    private final RuntimeService runtimeService;

    public WorkOrderProcessService(RuntimeService runtimeService) {
        this.runtimeService = runtimeService;
    }

    public void workOrderCompleted(String workOrderId, Map<String, Object> variables) {
        List<Execution> waiting = runtimeService.createExecutionQuery()
                .processDefinitionKey("work-order-execution")
                .messageEventSubscriptionName("WorkOrderCompleted")
                .variableValueEquals("workOrderId", workOrderId)
                .list();

        if (waiting.isEmpty()) {
            throw new NoWaitingInstanceException(workOrderId);
        }
        if (waiting.size() > 1) {
            throw new AmbiguousCorrelationException(workOrderId, waiting.size());
        }

        runtimeService.messageEventReceived(
                "WorkOrderCompleted", waiting.get(0).getId(), variables);
    }
}

Detrás hay una regla que en BPMN no aparece en negrita en ninguna parte y aun así lo determina todo: para un mensaje debe haber exactamente un punto de espera registrado. Una Message se dirige a un destinatario, y el motor tiene que entregarla exactamente a una Execution. Si no hay ninguna, no tiene a nadie. Si hay dos, no tiene una opción correcta, y entonces prefiere no elegir ninguna.

Esta llamada puede rechazar de cuatro maneras, y ninguna de ellas es un bug del motor:

  1. Sin resultado, porque la instancia todavía no llegó al Catch Event. La event subscription contra la que se correlaciona es una fila en ACT_RU_EVENT_SUBSCR, y solo nace cuando el token llega al Catch Event. Flowable no conoce un mecanismo de búfer para mensajes BPMN; se correlaciona exclusivamente contra event subscriptions existentes. Si el aviso de finalización del proveedor adelanta a tu propio proceso, por ejemplo porque el paso de entrega todavía está en marcha, la query no encuentra nada.
  2. Sin resultado, porque la instancia ya terminó. El mensaje llegó duplicado, o alguien empujó el proceso a mano. Desde fuera este caso no se distingue del primero: la query devuelve en ambos una lista vacía.
  3. Más de un resultado. Si dos Executions con la misma workOrderId esperan el mismo Message, list() devuelve ambas y el llamador tiene que decidir qué significa eso. Quien usa singleResult() en su lugar recibe la decisión como una FlowableException con el texto literal Query return 2 results instead of max 1, pero pierde la información de cuántas eran. Cómo llegan a existir dos Executions en espera lo muestra el capítulo sobre la reparación.
  4. La Execution existe, pero ya no tiene la subscription. Entre la query y la entrega hay una ventana de tiempo, y si la instancia avanza justo en ella, messageEventReceived lanza, según el Javadoc, una FlowableObjectNotFoundException si falta la Execution, o una FlowableException si no está suscrita al Message.

Los cuatro rechazos son afirmaciones sobre el estado del proceso, no sobre la infraestructura. Solo el primero se cura solo si se le da tiempo. Los otros tres persisten, da igual cuántas veces se repita la misma llamada. Esta distinción vuelve en cada variante.

El primer error: retry contra un rechazo de negocio

El primer reflejo contra el rechazo es un retry en el receptor: si la correlación falla, esperar un poco y volver a intentarlo.

private void process(BasicAcknowledgeablePubsubMessage message) throws Exception {
    WorkOrderCompleted event = read(message);

    for (int attempt = 1; attempt <= 30; attempt++) {
        try {
            processService.workOrderCompleted(event.workOrderId(), variables(event));
            message.ack();
            return;
        } catch (NoWaitingInstanceException e) {
            Thread.sleep(10_000);   // seguro que la instancia llega enseguida al Catch Event
        }
    }
}

Para el rechazo número uno, el de la instancia que todavía no está lista, esto incluso funciona: el tiempo cura ese caso. El lugar de la espera sigue siendo el equivocado, por dos razones.

La primera es el Head-of-Line-Blocking. El callback ocupa un thread de procesamiento, y el flow control de la client library limita cuántos mensajes pueden estar pendientes a la vez. Si el thread se pasa cinco minutos en el bucle de la orden 4711, las notificaciones de todas las demás órdenes esperan detrás de él. El Head-of-Line-Blocking es aquí una propiedad de la configuración del consumer, no del broker. Pub/Sub habría entregado los otros mensajes hace tiempo; es tu receptor el que no los acepta.

La segunda razón es la lease. La ack deadline llega como máximo a 600 segundos; todo lo que pase de ahí lo sostiene solo la client library, extendiendo la deadline en segundo plano hasta agotar su maxAckExtensionPeriod. Un retry que se pelea cien minutos contra un error acaba trabajando con un mensaje cuya lease puede haber caído hace tiempo. Pub/Sub ya lo habrá entregado de nuevo, quizá a otra instancia del work-order-service, y si tu ack tardío todavía surte algún efecto ya no lo sabe nadie.

Contra los rechazos dos a cuatro el retry no consigue nada de todos modos. Una instancia de proceso que ya terminó no vuelve por repetir la llamada, y un token duplicado no desaparece por eso. El retry convierte un rechazo de negocio en un bucle infinito con tiempo de espera.

El segundo error: el ack en el bloque catch

Tras la primera tormenta de retries que atascó la subscription llega, con puntualidad, la segunda reparación: en caso de error se hace ack para que la cosa siga. Exactamente así nació el bloque catch del receptor del principio.

Ahora hay que separar los dos defectos. Quien hace nack ante una excepción tiene un problema de disponibilidad y de acoplamiento: la subscription acumula mensajes mientras el motor no esté accesible. Pero el mensaje sigue existiendo, y cuando la avería pasa, vuelve. Quien hace ack ante una excepción convierte la misma avería en pérdida de datos: Pub/Sub borra el mensaje, sin «Retain acknowledged messages» ni topic retention no hay vuelta atrás, y el proceso de la orden espera para siempre.

Lo traicionero del asunto: este antipatrón no necesita un bloque catch escrito a mano. Surge también solo por configuración. Quien en vez de PubSubTemplate.subscribe va por el camino de Spring Integration con el PubSubInboundChannelAdapter y @ServiceActivator elige un AckMode, y sus variantes se comportan de forma radicalmente distinta en caso de error:

AckMode con éxito con excepción
AUTO ack sin error handler nack con reentrega inmediata, con un error handler que termina bien ack
AUTO_ACK ack sin error handler ninguna acción, la deadline expira; con error handler se hace ack
MANUAL nada nada, tu código hace ack o nack por sí mismo

La línea que arma la pérdida de datos parece del todo inofensiva:

@Bean
public PubSubInboundChannelAdapter workOrderAdapter(PubSubTemplate pubSubTemplate,
                                                    MessageChannel workOrderChannel) {
    var adapter = new PubSubInboundChannelAdapter(pubSubTemplate, "work-order-completed-sub");
    adapter.setOutputChannel(workOrderChannel);
    adapter.setAckMode(AckMode.AUTO_ACK);
    adapter.setErrorChannelName("pubsubErrors");   // a partir de aquí se hace ack en caso de error
    return adapter;
}

Con AUTO_ACK y un error handler, el mensaje se da por resuelto en cuanto el handler lo ha visto, aunque no haga nada más que loguear. Quien quiere control total usa MANUAL y saca la BasicAcknowledgeablePubsubMessage del header GcpPubSubHeaders.ORIGINAL_MESSAGE. Con esto la recomendación queda fijada antes de escribir una sola línea de código de inbox: un ack solo puede ir detrás del punto en el que el mensaje está a salvo de forma duradera.

Por qué la reparación construye la siguiente avería

Queda la pregunta de dónde sale el rechazo número tres, el token duplicado. No nace de un bug, sino de una reparación manual, y entre causa y efecto pueden pasar meses.

El antecedente: una notificación se perdió, digamos que por el ack en el bloque catch. La orden está parada. Alguien de operaciones interviene y, mediante una operación de change state, le coloca a la instancia un token en el paso de entrega para que la orden se entregue de nuevo.

El paso de entrega está modelado como asíncrono con flowable:async="true", como corresponde a una llamada externa. Y justo eso hace engañosa la intervención: la operación solo crea un job en ACT_RU_JOB; se ejecuta cuando el Job Executor lo recoge. Justo después de la intervención, por tanto, todo parece igual que antes. Ninguna llamada nueva en los logs, ningún avance visible en el diagrama del proceso. La conclusión obvia es: la intervención no funcionó, así que otra vez. Cada uno de esos intentos coloca un token más al lado.

En el momento de nacer, el token extra no causa ni un solo síntoma. Ambos tokens pasan por la entrega, ambos llegan al Catch Event, ambos crean una event subscription con la misma workOrderId, y el proceso se ve en el monitor como siempre. Solo cuando llega la siguiente notificación de esa orden, singleResult() falla con Query return 2 results instead of max 1. Si encima sigue corriendo el receptor con el ack en el catch, ese rechazo de negocio se loguea, el mensaje se borra y la orden vuelve a quedarse parada. La reparación de la última avería construyó la siguiente.

La lección es incómoda: un error de correlación puede deberse a una intervención de hace meses. Un sistema que tira esos mensajes en vez de conservarlos se quita a sí mismo toda opción de diagnóstico y reparación.

Variante 1: nack, retry policy y Dead-Letter-Topic

La primera variante sólida se queda por completo en Pub/Sub y no necesita infraestructura nueva: ante errores se hace nack, la repetición la asume una retry policy, y lo que fracasa de forma duradera va a parar a un Dead-Letter-Topic.

El receptor se vuelve más corto, no más largo:

private void process(BasicAcknowledgeablePubsubMessage message) {
    try {
        WorkOrderCompleted event = read(message);
        processService.workOrderCompleted(event.workOrderId(), variables(event));
        message.ack();
    } catch (Exception e) {
        log.warn("Notificación aplazada: {}", e.getMessage());
        message.nack();
    }
}

Dos detalles de configuración deciden si esto acaba bien, y ambos se pasan por alto con regularidad.

Primero, la retry policy. Sin policy configurada, el default de Pub/Sub es la reentrega inmediata, sin ningún backoff. Un nack a un rechazo de negocio genera entonces un bucle denso de entrega, rechazo, entrega. Con policy rige backoff exponencial, minimumBackoff default 10 segundos, maximumBackoff default 600 segundos, ambos ajustables entre 0 y 600 segundos. La policy se aplica por mensaje, tanto con nack como al expirar la ack deadline.

Segundo, el Dead-Letter-Topic. Se configura en la subscription, con maxDeliveryAttempts entre 5 y 100, default 5:

gcloud pubsub topics create work-order-completed-dlt

gcloud pubsub subscriptions update work-order-completed-sub 
    --dead-letter-topic=work-order-completed-dlt 
    --max-delivery-attempts=10 
    --min-retry-delay=10s 
    --max-retry-delay=600s

Esto necesita dos IAM bindings sin los cuales sencillamente no hay reenvío: la service account de Pub/Sub service-<número de proyecto>@gcp-sa-pubsub.iam.gserviceaccount.com necesita roles/pubsub.publisher en el Dead-Letter-Topic y roles/pubsub.subscriber en la subscription de origen.

gcloud pubsub topics add-iam-policy-binding work-order-completed-dlt 
    --member="serviceAccount:service-123456789@gcp-sa-pubsub.iam.gserviceaccount.com" 
    --role="roles/pubsub.publisher"

gcloud pubsub subscriptions add-iam-policy-binding work-order-completed-sub 
    --member="serviceAccount:service-123456789@gcp-sa-pubsub.iam.gserviceaccount.com" 
    --role="roles/pubsub.subscriber"

La honestidad exige nombrar los límites de la promesa: el reenvío al Dead-Letter-Topic es best-effort; pueden producirse menos o más intentos de entrega que los configurados. El mensaje reenviado va envuelto y lleva atributos CloudPubSubDeadLetterSource, entre ellos la subscription de origen y el contador de entregas. El campo delivery_attempt, en cambio, es lo que ven los subscribers de la subscription de origen en cada entrega. El contador detrás de ese atributo solo se lleva si el Dead-Letter-Topic está configurado correctamente, y puede volver a caer a 0; el caso típico son subscriptions pull cuyos subscribers están inactivos a ratos.

Muchos artículos sobre la inbox callan este punto: para un consumer sin base de datos propia, con volumen moderado y sin requisitos de orden, esta variante es en muchos casos la mejor solución. Ningún code path nuevo, ninguna tabla, ningún relay, y la monitorización viene de serie, por ejemplo con la métrica oldest_unacked_message_age y la ocupación del Dead-Letter-Topic. Lo que le falta solo se ve en cuatro requisitos concretos.

Variante 2: la Transactional Inbox

La Transactional Inbox, la contraparte del lado receptor de la Transactional Outbox de Chris Richardson, separa dos cosas que el receptor del principio mezcló: la recepción del mensaje y su procesamiento. El receptor solo escribe el mensaje en una tabla propia y hace ack. La correlación con el motor la asume más tarde un relay separado.

La tabla:

CREATE TABLE work_order_inbox (
    id             BIGINT GENERATED ALWAYS AS IDENTITY PRIMARY KEY,
    event_type     VARCHAR(64)  NOT NULL,
    work_order_id  VARCHAR(36)  NOT NULL,
    payload        JSONB        NOT NULL,
    status         VARCHAR(16)  NOT NULL DEFAULT 'OPEN',
    attempts       INT          NOT NULL DEFAULT 0,
    next_attempt   TIMESTAMPTZ  NOT NULL DEFAULT now(),
    last_error     TEXT,
    received_at    TIMESTAMPTZ  NOT NULL DEFAULT now(),
    processed_at   TIMESTAMPTZ,
    CONSTRAINT uq_work_order_inbox UNIQUE (event_type, work_order_id)
);

El receptor se reduce a recibir, guardar, confirmar:

private void process(BasicAcknowledgeablePubsubMessage message) {
    try {
        WorkOrderCompleted event = read(message);
        inbox.store("WorkOrderCompleted", event.workOrderId(), rawData(message));
        message.ack();
    } catch (Exception e) {
        log.warn("No se pudo guardar la notificación: {}", e.getMessage());
        message.nack();
    }
}
@Repository
public class WorkOrderInbox {

    private final JdbcTemplate jdbc;

    public WorkOrderInbox(JdbcTemplate jdbc) {
        this.jdbc = jdbc;
    }

    public boolean store(String eventType, String workOrderId, String payload) {
        int inserted = jdbc.update("""
                INSERT INTO work_order_inbox (event_type, work_order_id, payload)
                VALUES (?, ?, ?::jsonb)
                ON CONFLICT ON CONSTRAINT uq_work_order_inbox DO NOTHING
                """, eventType, workOrderId, payload);
        return inserted == 1;
    }
}

En este código pequeño hay tres decisiones.

Entre el insert y el ack queda una ventana de crash. Si el pod muere justo entre las dos líneas, el mensaje está en la tabla pero sin confirmar, y Pub/Sub lo entrega de nuevo. El insert tiene que sobrevivir sin daños a esa repetición; por eso es idempotente a través del unique key: el segundo insert del mismo mensaje choca contra el constraint, no hace nada, y el ack sale igualmente. Eso es exactamente el Idempotent Consumer de microservices.io, vertido en una tabla.

La elección del unique key no es trivial. (event_type, work_order_id) solo funciona mientras el proveedor mande como máximo un aviso de finalización por orden. Si dos mensajes distintos de negocio llevan el mismo key, por ejemplo porque una orden se cierra en dos tramos, el unique key se traga el segundo sin comentario, sin error y sin línea de log. Entonces el key necesita un ID de evento de negocio del emisor, y si el emisor no lo entrega, eso es una conversación con el emisor, no un detalle de implementación.

Y sí: la base de datos está ahora ella misma en la ruta del ack. Eso no contradice la crítica al receptor del principio, sino que la resuelve. Antes del ack puede estar lo que solo puede fallar de forma técnica, porque un fallo técnico lo cura la reentrega de Pub/Sub. Un INSERT en una tabla propia no puede rechazar por motivos de negocio; la correlación con el motor sí puede. Por eso el insert va antes del ack y la correlación detrás.

El relay, y la trampa del orden

La segunda pieza es el relay: un poller que lee las entradas abiertas de la inbox y las correlaciona con el motor. La forma obvia es esta:

@Component
public class WorkOrderInboxRelay {

    private final JdbcTemplate jdbc;
    private final WorkOrderProcessService processService;

    public WorkOrderInboxRelay(JdbcTemplate jdbc, WorkOrderProcessService processService) {
        this.jdbc = jdbc;
        this.processService = processService;
    }

    @Scheduled(fixedDelay = 2000)
    @Transactional
    public void processOpenEntries() {
        List<WorkOrderInboxEntry> entries = jdbc.query("""
                SELECT * FROM work_order_inbox
                WHERE status = 'OPEN' AND next_attempt <= now()
                ORDER BY id
                LIMIT 10
                FOR UPDATE SKIP LOCKED
                """, entryMapper());

        for (WorkOrderInboxEntry entry : entries) {
            process(entry);
        }
    }
}

FOR UPDATE SKIP LOCKED hace el relay escalable en horizontal: varios pods del work-order-service hacen polling a la vez, cada uno bloquea sus filas, ninguno espera al otro. Es el patrón Competing Consumers de Hohpe y Woolf, trasladado a una tabla.

Y justo ahí está la trampa. En cuanto la inbox transporta más de un tipo de evento por orden, por ejemplo «ejecución iniciada» para el historial de estado y «orden completada» para la correlación, dos pods pueden tomar a la vez los dos mensajes de la misma orden. Entonces quizá se correlaciona «completada» antes de que «iniciada» esté procesada, el proceso ya pasó de largo del segundo evento, y la correlación de «iniciada» rechaza. Es exactamente el race del rechazo número uno, solo que reconstruido un piso más abajo, en tu propia infraestructura en vez de en el broker.

La solución consiste en serializar por Correlation Key y seguir en paralelo entre keys. En SQL se puede expresar directamente: una entrada solo puede tomarse si para la misma orden no existe una entrada más antigua sin resolver.

SELECT i.*
FROM work_order_inbox i
WHERE i.status = 'OPEN'
  AND i.next_attempt <= now()
  AND NOT EXISTS (
      SELECT 1
      FROM work_order_inbox older
      WHERE older.work_order_id = i.work_order_id
        AND older.status IN ('OPEN', 'IN_PROGRESS')
        AND older.id < i.id)
ORDER BY i.id
LIMIT 10
FOR UPDATE SKIP LOCKED;

Si el pod A bloqueó la fila más antigua de una orden, el pod B la sigue viendo como pendiente y deja la más nueva en su sitio. En el caso límite una entrada espera con esto un ciclo de polling más de lo necesario; esa es la dirección segura. Quien quiera ahorrarse la construcción tiene una opción más simple con un precio claro: exactamente una instancia de relay, ORDER BY id estricto, y el throughput queda limitado a ese único worker. Para un puñado de órdenes por hora es perfectamente aceptable.

¿Error técnico o de negocio?

En el relay se toma ahora la decisión que el receptor original nunca tomó: ¿qué significa un error y qué se sigue de él? La clasificación habitual en sistemas distribuidos separa errores transitorios, que se repiten con backoff exponencial, de errores permanentes, que van directos a la estación final. Para la correlación esto significa en concreto:

private void process(WorkOrderInboxEntry entry) {
    try {
        processService.workOrderCompleted(entry.workOrderId(), entry.variables());
        markCompleted(entry);
    } catch (NoWaitingInstanceException e) {
        // quizá se cura con el tiempo: backoff, pero con plazo
        if (entry.attempts() >= maxAttempts) {
            park(entry, e);
        } else {
            rescheduleWithBackoff(entry, e);
        }
    } catch (AmbiguousCorrelationException e) {
        // dos puntos de espera: no hay retry que ayude
        park(entry, e);
    } catch (FlowableException e) {
        // la subscription desapareció entre la query y la entrega
        park(entry, e);
    } catch (Exception e) {
        // técnico: base de datos, motor no accesible, timeout
        rescheduleWithBackoff(entry, e);
    }
}

La asignación se deriva del capítulo de los cuatro rechazos. Ningún resultado puede significar que la instancia todavía va de camino, así que la paciencia compensa, pero no sin límite: pasado un plazo, un «todavía no está» se convirtió en un «nunca estuvo o ya terminó», y la entrada se aparca. Un resultado múltiple y una subscription desaparecida son de inmediato un caso para personas; cada reintento automático solo produciría el mismo rechazo. Y todo lo técnico recibe backoff, porque ahí la nueva ejecución sí cura.

Queda una imprecisión, y conviene conocerla: FlowableException no es una excepción puramente de negocio. El motor envuelve en ella también errores técnicos, por ejemplo cuando la base de datos se le cae debajo. Quien quiera separar el caso con limpieza no se libra de examinar la causa en vez de fiarse del tipo de la excepción. Las dos excepciones del propio servicio, NoWaitingInstanceException y AmbiguousCorrelationException, son tipos propios precisamente por eso, y no errores del motor pasados tal cual.

Aparcada significa: status = 'PARKED', el último error queda en la fila, y salta una alarma. Operaciones necesita como mínimo dos alarmas, una sobre entradas aparcadas y otra sobre la edad de la entrada abierta más antigua. Ambas son queries SQL sencillas, y que haya que construirlas y operarlas uno mismo no es un detalle, sino un coste real de esta arquitectura. El capítulo sobre los límites de la inbox vuelve sobre ello.

Variante 3: el motor consume por sí mismo

Si la lógica de correlación cuelga de todos modos tan cerca del motor, la idea es tentadora: ¿por qué no consume el propio motor del broker? Algunos motores traen para ello una integración de eventos en la que una subscription del broker se liga declarativamente al modelo de proceso y el receiver escrito a mano desaparece por completo.

El atractivo es real: una capa menos, ninguna lógica de correlación propia, el mapeo de mensaje a Correlation Key está en el modelo en vez de en el código Java. Para casos simples es la solución más ligera.

Tres cosas conviene ver con realismo. Primero, el rechazo de correlación no desaparece, solo se muda: también un adapter del motor se topa con instancias que todavía no esperan o ya no esperan, y cómo maneja entonces ack, retry y estación final lo determina ahora la implementación del adapter en vez de tu código. Las preguntas de este artículo hay que hacérselas igualmente al adapter, solo que allí las responde otro. Segundo, el consumer del broker entra en el ciclo de vida del motor: un reinicio del motor es ahora automáticamente también un reinicio del consumer; el acoplamiento que queríamos sacar de la ruta del ack vuelve en otra forma. Y tercero, la integración es cuestión de lo que el motor trae de serie: qué broker se soporta depende del motor y de la versión, y para Google Cloud Pub/Sub, según el stack, la cosa acaba en un adapter mantenido por ti, con lo que el código supuestamente ahorrado vuelve a estar ahí.

Como regla rápida sirve esto: la variante despliega su fuerza cuando el adapter para tu broker existe, su manejo de errores está documentado y es aceptable, y no tienes requisitos propios de deduplicación o reparación. Si no, el camino explícito vía receiver o inbox sigue siendo la elección más controlable.

Variante 4: un motor que almacena mensajes en búfer

Los cuatro rechazos del capítulo de correlación tienen una raíz común: Flowable correlaciona exclusivamente contra event subscriptions que existen en el momento de la llamada. Hay motores construidos de otra manera en este punto: los mensajes publicados se almacenan en búfer en el broker del motor, con un time to live cuyo parámetro se llama timeToLive y se indica en milisegundos. Con la TTL a 0, el búfer está apagado. Mientras la TTL corre, un mensaje en búfer correlaciona incluso si la subscription adecuada nace después de la publicación. En una implementación extendida el cliente pone por defecto una hora, y hay además una messageId opcional para idempotencia.

Trasladado a nuestro escenario, el rechazo número uno desaparece por completo: la notificación puede adelantar al proceso, simplemente espera en el motor hasta que el token llega al Catch Event. El race entre entrega y notificación, contra el que se construyeron el retry en el receptor y el plazo en el relay, sencillamente no existe en un modelo así.

Hay que ser honesto con los demás rechazos. A una instancia que ya terminó tampoco la recupera un búfer una vez expirada la TTL. Un token duplicado sigue siendo un token duplicado. Quien usa la llamada de correlación síncrona separada de esos motores no recibe además ningún búfer; ese vale solo para la vía de publicación. Y para Signals no vale en absoluto. El búfer quita de la mesa uno de los cuatro rechazos; los otros tres se quedan.

Para quien apuesta por Flowable, la consecuencia es más simple: el búfer que otros tienen en el broker es aquí la tabla de inbox. El mismo concepto, solo que operado por ti y, a cambio, accesible por SQL.

Variante 5: Claim Check

Claim Check viene de los Enterprise Integration Patterns de Hohpe y Woolf. El patrón resuelve en realidad otro problema: mensajes grandes. El payload va a un almacén de datos, por el broker fluye solo la referencia, y el receptor recoge los datos al procesar.

En nuestro escenario habría incluso un motivo real para ello: el protocolo de cierre completo de una orden con varios miles de posiciones no pinta nada en un mensaje de Pub/Sub. El proveedor deja el protocolo en un bucket, el mensaje lleva workOrderId y la referencia, y la correlación trabaja solo con el mensaje ligero.

Al problema de correlación, sin embargo, Claim Check solo contribuye de forma indirecta, y no conviene venderlo como algo más grande de lo que es. La contribución indirecta: los datos de negocio quedan de forma duradera en el store, con independencia de ack, retention y conservación en el dead letter. Incluso si el mensaje se pierde definitivamente, del store se puede reconstruir qué órdenes se completaron, y es posible rehacer la correlación a mano o por script. Es una red de seguridad debajo de la red de seguridad, pero no correlaciona nada por sí sola: los rechazos del motor, la ventana de crash en el receptor y el orden en el relay quedan exactamente como estaban.

Lo que la inbox no resuelve

Un capítulo propio para los límites, para que no desaparezcan en la letra pequeña.

La inbox no cura la causa. El token duplicado del capítulo de la reparación sigue siendo un token duplicado; la correlación sigue rechazando su notificación, venga del receiver o del relay. Lo que la inbox cambia es el desenlace: de una prueba borrada sale una fila aparcada con payload y texto de error. El fallo se vuelve sobrevivible y reparable, no imposible. La reparación en sí es entonces deliberadamente poco espectacular:

UPDATE work_order_inbox
SET status = 'OPEN', attempts = 0, next_attempt = now()
WHERE id = 4711;

Primero eliminar el token duplicado en el motor, luego esta única línea, y el relay hace el resto. Sin tooling de re-publish, sin copiar del Dead-Letter-Topic.

A cambio te compras costes de operación, y el reproche que suele caer aquí está justificado: estás construyendo una queue delante de la queue. La acumulación se muda de una subscription, que trae métricas hechas como oldest_unacked_message_age y alerting rodado, a una tabla que al principio no vigila nadie. A eso se suman el crecimiento y el cleanup de las filas resueltas, un segundo mecanismo de retry junto a la retry policy de Pub/Sub y una segunda estación final junto al Dead-Letter-Topic. Cada una de esas piezas hay que construirla, probarla y entenderla en el equipo.

¿Cuándo compensa aun así? La comparación con la variante 1 se puede condensar en cuatro requisitos. El reloj de retention: también en el Dead-Letter-Topic sigue corriendo; a más tardar tras 31 días el mensaje desaparece también allí, las filas de la inbox no caducan nunca. Deduplicación: el unique key de la inbox hace al consumer idempotente; Pub/Sub por sí solo no lo hace. Reparación por SQL en vez de tooling de re-publish, como se mostró arriba. Y una transacción local existente en la que el insert deba entrar de forma atómica, por ejemplo cuando la recepción escribe de todos modos datos de negocio. Si no se cumple ninguno de los cuatro, nack con retry policy y Dead-Letter-Topic es la elección más simple y, por eso, la mejor.

El worker reanudable

Con la correlación la historia no termina, porque detrás del Catch Event el proceso continúa, y allí rigen las mismas leyes. El paso «emitir la factura» está modelado con flowable:async="true", y lo que eso significa en una caída conviene haberlo visto una vez en detalle.

flowable:async="true" crea un job en ACT_RU_JOB. Cuando un Async Executor recoge el job, escribe lock owner y lock expiration en la fila. Si el pod muere en plena ejecución, el lock se queda hasta que expira, entonces el reset thread lo libera y otro executor ejecuta el job. Y lo hace desde el principio, por completo: lo que el pod muerto ya hizo no lo sabe nadie. Eso es at-least-once, la misma semántica que en Pub/Sub, solo que un nivel más arriba, y la consecuencia es la misma: la facturación tiene que ser idempotente, o la misma orden se factura dos veces en la repetición.

Los defaults, todos configurables:

Ajuste Default
asyncExecutorNumberOfRetries 3, después tabla de deadletter jobs
asyncExecutorAsyncJobLockTimeInMillis 5 minutos
asyncExecutorTimerLockTimeInMillis 5 minutos
asyncExecutorResetExpiredJobsInterval 60 segundos
asyncExecutorResetExpiredJobsPageSize 3
asyncExecutorDefaultAsyncJobAcquireWaitTime 10 segundos

Dos de esos números son trampas. Tres retries se agotan rápido si el sistema de facturación se atasca diez minutos; después el job está en la tabla de deadletter y vuelve a necesitar a una persona. Y cinco minutos de lock significan: tras la muerte de un pod, el paso se queda parado hasta cinco minutos más el intervalo de reset antes de que alguien lo retome.

Quien quiera sacar la facturación del motor usa el External Worker: el proceso deja un job en un topic, y un worker externo se lo lleva con un derecho exclusivo limitado en el tiempo, una lease.

String workerId = "billing-worker-1";

List<AcquiredExternalWorkerJob> jobs = managementService
        .createExternalWorkerJobAcquireBuilder()
        .topic("create-invoice", Duration.ofMinutes(10L))
        .acquireAndLock(20, workerId);

for (AcquiredExternalWorkerJob job : jobs) {
    try {
        createInvoice(job);
        managementService
                .createExternalWorkerCompletionBuilder(job.getId(), workerId)
                .complete();
    } catch (Exception e) {
        managementService
                .createExternalWorkerJobFailureBuilder(job.getId(), workerId)
                .errorMessage(e.getMessage())
                .fail();
    }
}

Si el procesamiento fracasa, el worker comunica el fallo en vez de complete(), y el job se vuelve a repartir. Una extensión de la lease en curso no está documentada en el cliente External de Java; por eso la duración del lock en la llamada a topic() debe quedar holgadamente por encima de la duración de procesamiento esperada: si la lease expira en mitad del trabajo, un segundo worker puede llevarse el mismo job, y la idempotencia tiene que soportar también ese caso.

¿Signal o Message?

Al modelar la notificación surge con regularidad la pregunta de si un Signal no sería más simple que la Message, justo porque los Signals no lanzan error cuando falta el receptor. La respuesta para este escenario es inequívoca.

En BPMN 2.0 un Signal solo tiene un nombre y ningún destinatario; una Message tiene emisor y receptor. Un Signal actúa en global: alcanza a toda instancia en marcha que justo lo esté esperando, y ninguna de ellas está nombrada en el Signal. signalEventReceived(String signalName) notifica en consecuencia a todas las Executions con signal subscription activa. Si no existe ni una sola, no pasa nada: sin excepción, sin búfer, el Signal se desvanece. A eso se añade que el lanzamiento del Signal es síncrono por defecto; el proceso que lo lanza espera, por tanto, hasta que el Signal se ha entregado a todos los catchers.

Exactamente la propiedad que hace parecer cómodo al Signal lo descalifica para el aviso de finalización. «Orden completada» para la orden 4711 es una afirmación dirigida a exactamente una instancia de proceso, y si esa instancia no espera, eso es una información que el emisor u operaciones tienen que conocer. La Message entrega esa información como rechazo, con el que los capítulos de arriba saben lidiar. El Signal no la entrega: reporta el mismo estado como éxito sin efecto, y la pérdida ruidosa vuelve a ser silenciosa. Un Signal encaja donde de verdad muchas instancias tienen interés en la misma novedad, por ejemplo un cambio de tarifa del proveedor que afecta a todas las órdenes en marcha. Una notificación dirigida con Correlation Key es una Message, en cada uno de los modelos mostrados aquí. Qué pasa si aun así te decides por el broadcast, y cómo hacer visible el daño en un test, está en BPMN Signal o Message.

Qué variante y cuándo

La tabla de decisión, con el receptor del principio como vara de medir de lo que hay que evitar:

Variante Úsala si Déjala si
Nack, retry policy, Dead-Letter-Topic consumer sin base de datos propia, volumen moderado, sin requisito de orden, el monitoring estándar basta los mensajes tienen que sobrevivir más que la retention (máx. 31 días, también en el DLT), se necesita deduplicación o se quiere reparación por SQL
Transactional Inbox con relay cuenta al menos uno: reloj de retention, deduplicación vía unique key, reparación por SQL, o un insert que pertenece de forma atómica a una transacción local existente no se cumple ninguno de los cuatro criterios, porque entonces pagas tabla, relay, monitoring y cleanup sin contrapartida
El motor consume por sí mismo existe un adapter terminado y documentado para tu broker y su manejo de errores cumple tus requisitos tienes requisitos propios de deduplicación, backoff o estación final, o quieres desplegar el consumer con independencia del motor
Motor con búfer de mensajes la elección de motor sigue abierta y el race de adelantamiento es tu problema principal el motor ya está fijado; en Flowable la inbox asume el papel de búfer
Claim Check el payload es grande y debe nacer un fondo de datos de negocio duradero junto al messaging esperas de él una solución del problema de correlación, porque no la entrega

Por encima de todas las variantes hay dos reglas que valen siempre. Primera: un ack solo cuando el mensaje está a salvo de forma duradera, y antes del ack solo pasos que únicamente pueden fallar de forma técnica. Segunda: los rechazos de negocio de la correlación necesitan una estación final y una alarma, no un retry.

FAQ

¿Por qué mi mensaje de Pub/Sub desapareció aunque el código lanzó una excepción?
Porque en algún sitio se hizo ack de todos modos. Los dos caminos habituales: un ack() en el bloque catch, o AckMode.AUTO_ACK en combinación con un error handler, porque esa combinación confirma el mensaje también en caso de error. Sin «Retain acknowledged messages» en la subscription o message retention en el topic, un mensaje con ack no es recuperable.

¿Resuelve el problema exactly-once delivery de Pub/Sub?
No. Exactly-once delivery existe solo para subscriptions pull, vale solo dentro de una región y garantiza exactly-once delivery, no exactly-once processing. Los duplicados del lado del publisher siguen además siendo posibles. El rechazo de negocio de la correlación, el problema central de este artículo, ni lo toca.

¿Por qué Flowable no encuentra mi instancia de proceso aunque está corriendo?
Porque se correlaciona antes de que la instancia haya llegado al Catch Event. La event subscription nace cuando el token llega al Catch Event, y Flowable no conoce un mecanismo de búfer para mensajes BPMN. Una notificación que adelanta a tu propio proceso no encuentra por eso ninguna subscription y hay que reintentarla más tarde.

¿Ayudan las ordering keys contra la trampa del orden en el relay?
Solo del lado del broker, y con etiqueta de precio: por ordering key es posible 1 MBps de throughput, y la redelivery de un mensaje arrastra consigo todos los mensajes posteriores de la misma key, también los ya confirmados. En cuanto un relay paralelo tira de una tabla, la serialización por Correlation Key tiene que ocurrir allí de todos modos; las ordering keys de la subscription ya no ayudan entonces.

¿Sigo necesitando la inbox si ya tengo un Dead-Letter-Topic?
A menudo no. El Dead-Letter-Topic ya cubre el caso «fracasó de forma duradera, tiene que intervenir una persona». La inbox solo gana cuando cuenta uno de cuatro requisitos: conservación más allá de la retention, porque también en el Dead-Letter-Topic rigen 31 días como máximo, deduplicación vía unique key, reparación por SQL, o un insert atómico en una transacción local existente.

¿Cómo evito que la tabla de inbox crezca sin límite?
Con un cleanup job propio que nadie más te da: borrar o archivar las filas resueltas tras un plazo de conservación, excluyendo expresamente las filas aparcadas. A eso pertenecen dos alarmas, una sobre entradas aparcadas y otra sobre la edad de la entrada abierta más antigua, porque las métricas hechas de Pub/Sub no ven la tabla.

¿Qué pasa con un mensaje que en 31 días no se procesó en ningún sitio?
Desapareció. La message retention de una subscription es de 31 días como máximo, y tras su expiración Pub/Sub puede descartar el mensaje con independencia del estado de confirmación. Eso vale también para la subscription de un Dead-Letter-Topic. Quien necesite conservar más tiempo necesita un almacén propio, por ejemplo la inbox o un store de Claim Check.

Conclusión

Antes del ack solo puede estar lo que únicamente puede fallar de forma técnica. Un fallo técnico lo cura la reentrega del broker; un rechazo de negocio de la correlación no lo cura ningún retry del mundo, y por eso el motor de workflow va detrás del ack, no delante.

En el camino hasta ahí hay dos defectos que conviene mantener separados. El motor en la ruta del ack es un problema de disponibilidad: cada deployment atasca la subscription, pero los mensajes sobreviven. El ack en caso de error, sea como bloque catch o como AUTO_ACK con error handler, lo convierte en pérdida de datos, y esa, en Pub/Sub y sin precauciones, es definitiva.

Para la solución rige este orden: comprobar primero nack con retry policy y Dead-Letter-Topic, porque para muchos consumers es la mejor elección, con métricas y alerting de serie. La Transactional Inbox entra en juego cuando cuentan retention, deduplicación, reparación por SQL o una transacción local, y llega con costes de operación palpables: vigilancia propia, cleanup propio, estación final propia, y la trampa del orden en el relay paralelo, que obliga a serializar por Correlation Key. Y sea cual sea la variante final: las causas de los rechazos de negocio, desde el race de adelantamiento hasta el token duplicado de la reparación manual, no desaparecen con ninguna de ellas. Visible, conservado y reparable: más que eso no puede dar una arquitectura de entrega.

Fuentes

Todos los ejemplos de código de este artículo son propios.

Estado: agosto de 2026, verificado contra Spring Cloud GCP 8.1 y Flowable 8.0.

$ lang DE EN ES