Sincronizar el read model en CQRS: outbox, idempotencia y rebuild
Un lunes por la mañana, soporte abrió un ticket con la frase que más miedo da de todas: «el cliente dice que pagó y su pedido sigue saliendo como pendiente».
Miramos la tabla de pedidos. Pagado. Miramos la vista que consume el frontend. Pendiente. Llevaba once días así.
Nadie se había enterado porque el sistema no estaba roto. Todos los endpoints devolvían 200. Los logs estaban limpios. Simplemente, el read model se había quedado atrás y no existía ninguna alarma que mirase esa diferencia.
Ese es el problema real de CQRS: cómo sincronizar el read model con el write model sin que se pudra. No aparece el día que lo eliges; aparece tres meses después. Si todavía estás decidiendo si el patrón te conviene, ese es otro debate — qué es CQRS y cuándo compensa aplicarlo. Este post empieza el día siguiente.
Y la tesis, por delante: no pierdes eventos por culpa del bus. Los pierdes en el hueco que hay entre tu COMMIT y tu publish.
Por qué el read model se desincroniza: el problema del dual write
El read model se desincroniza porque el cambio de estado y la publicación del evento son dos operaciones contra dos sistemas distintos, sin transacción común. Si la segunda falla, no queda nadie para reintentarla.
Este código lo he visto en producción más veces de las que me gustaría:
await db.query(`update orders set status = 'paid' where id = $1`, [orderId])
await bus.publish(39;order.paid39;, { orderId })
Dos líneas. Parecen una unidad. No lo son.
Entre la primera y la segunda cabe todo: un timeout del broker, un deploy que mata el pod, el OOM killer, un ECONNRESET. Si la segunda línea falla, la base de datos de escritura dice "pagado" y el read model no se entera nunca. No hay reintento que te salve, porque el proceso que tenía que reintentar ya no existe.
Invertir el orden es peor. Si publicas primero y la transacción hace rollback después, has emitido un evento sobre algo que no ocurrió. El read model muestra un pedido pagado que en la fuente de verdad sigue pendiente. Ese bug se tarda semanas en encontrar.
Esto tiene nombre: dual write. Escribir en dos sistemas que no comparten transacción. No se arregla con try/catch, ni con reintentos en el catch, ni metiendo el publish dentro de la transacción —porque la red no hace rollback.
La solución académica es 2PC. La que usa la gente que tiene que dormir por las noches es otra.
Transactional outbox: sincronizar el read model en una sola transacción
El transactional outbox es un patrón que elimina el dual write escribiendo el evento en una tabla de la misma base de datos, dentro de la misma transacción que el cambio de estado; un proceso aparte —el relay— lee esa tabla y lo publica en el bus. Está catalogado así en el catálogo de patrones de microservicios de Chris Richardson.
Y la idea es tonta de simple: si no puedes hacer atómicas dos operaciones contra dos sistemas, haz que las dos vayan contra el mismo sistema.
El evento no se publica. Se inserta en una tabla de la misma base de datos, dentro de la misma transacción que el cambio de estado. Si el COMMIT pasa, el evento existe. Si no pasa, tampoco. Atomicidad gratis, la que ya te da Postgres.
create table outbox (
id bigserial primary key,
aggregate_id uuid not null,
aggregate_type text not null,
event_type text not null,
version int not null,
payload jsonb not null,
occurred_at timestamptz not null default now(),
published_at timestamptz
);
-- índice parcial: solo lo pendiente, que es lo que el relay consulta cada tick
create index outbox_pending_idx on outbox (id) where published_at is null;
Y el comando queda así:
import { Pool } from 39;pg39;
const pool = new Pool({ connectionString: process.env.DATABASE_URL })
export async function markOrderAsPaid(orderId: string, total: number) {
const client = await pool.connect()
try {
await client.query(39;begin39;)
const { rows } = await client.query<{ version: number }>(
`update orders
set status = 'paid', paid_at = now(), version = version + 1
where id = $1 and status = 'pending'
returning version`,
[orderId],
)
if (rows.length === 0) throw new Error(39;ORDER_NOT_PENDING39;)
await client.query(
`insert into outbox (aggregate_id, aggregate_type, event_type, version, payload)
values ($1, 'order', 'order.paid', $2, $3)`,
[orderId, rows[0].version, { orderId, total, paidAt: new Date().toISOString() }],
)
await client.query(39;commit39;)
} catch (error) {
await client.query(39;rollback39;)
throw error
} finally {
client.release()
}
}
Fíjate en el returning version. Ese número lo vas a necesitar dentro de dos secciones y es lo que separa un read model correcto de uno que a veces acierta.
El relay
Un proceso aparte lee la tabla y publica. Nada más. La clave está en el for update skip locked —disponible desde Postgres 9.5—: te deja correr varias instancias del relay sin que dos cojan la misma fila.
export async function relayTick(bus: Bus, batchSize = 100) {
const client = await pool.connect()
try {
await client.query(39;begin39;)
const { rows } = await client.query<OutboxRow>(
`select id, aggregate_id, event_type, version, payload, occurred_at
from outbox
where published_at is null
order by id
limit $1
for update skip locked`,
[batchSize],
)
for (const row of rows) {
await bus.publish({
id: String(row.id),
type: row.event_type,
aggregateId: row.aggregate_id,
version: row.version,
payload: row.payload,
occurredAt: row.occurred_at,
})
}
if (rows.length > 0) {
await client.query(
`update outbox set published_at = now() where id = any($1::bigint[])`,
[rows.map((r) => r.id)],
)
}
await client.query(39;commit39;)
} catch (error) {
await client.query(39;rollback39;)
throw error
} finally {
client.release()
}
}
Léelo otra vez y busca el agujero, porque lo tiene: si bus.publish va bien y el commit del update ... published_at falla, el evento sale publicado dos veces.
Eso es intencionado. El outbox te garantiza at-least-once, nunca exactly-once. Y está bien. Preferimos un evento duplicado que un evento perdido, porque el duplicado se resuelve en el consumidor y la pérdida no se resuelve en ningún sitio.
El relay es, además, donde el bus te va a fallar de verdad. Con el broker caído, reintentar en bucle cerrado solo empeora las cosas: aplica el mismo razonamiento que expliqué sobre circuit breakers y clasificación de fallos. Cortar, esperar, dejar que el outbox acumule. Para eso está la tabla: el relay puede pasarse diez minutos parado y no se pierde ni un evento.
Idempotencia en el proyector: procesar dos veces sin duplicar
Si el bus entrega al menos una vez, el proyector tiene que poder comerse el mismo evento dos veces y terminar en el mismo estado. Punto.
Dos piezas: una tabla que registra qué eventos ya se procesaron y un upsert que no dependa del orden de llegada.
create table processed_events (
projection text not null,
event_id bigint not null,
processed_at timestamptz not null default now(),
primary key (projection, event_id)
);
create table projection_checkpoint (
projection text primary key,
last_event_id bigint not null default 0,
updated_at timestamptz not null default now()
);
create table orders_read (
order_id uuid primary key,
status text not null,
total numeric not null,
paid_at timestamptz,
version int not null
);
El primary key de order_id no es decorativo: sin esa restricción única, el on conflict (order_id) del proyector ni siquiera llega a ejecutarse.
Antes de tocar nada, valida el payload. El evento viaja como JSON opaco y puede llevar meses en la tabla: el día que alguien cambie su forma en el productor, tu proyector recibirá algo que no espera. Parsea siempre con un schema —aquí, Zod 4— y manda a dead letter lo que no cumpla, en vez de dejar que un undefined acabe escrito en la vista.
import { z } from 39;zod39;
const OrderPaid = z.object({
orderId: z.uuid(),
total: z.number().nonnegative(),
paidAt: z.iso.datetime(),
})
export async function projectOrderPaid(event: DomainEvent) {
const parsed = OrderPaid.safeParse(event.payload)
if (!parsed.success) {
await deadLetter(event, parsed.error)
return
}
const client = await pool.connect()
try {
await client.query(39;begin39;)
// 1. Reclamar el evento. Si ya estaba, no hacemos nada más.
const claim = await client.query(
`insert into processed_events (projection, event_id)
values ('orders_read', $1)
on conflict do nothing`,
[event.id],
)
if (claim.rowCount === 0) {
await client.query(39;rollback39;)
return
}
// 2. Upsert con guarda de versión.
await client.query(
`insert into orders_read (order_id, status, total, paid_at, version)
values ($1, 'paid', $2, $3, $4)
on conflict (order_id) do update
set status = excluded.status,
total = excluded.total,
paid_at = excluded.paid_at,
version = excluded.version
where orders_read.version < excluded.version`,
[parsed.data.orderId, parsed.data.total, parsed.data.paidAt, event.version],
)
// 3. Avanzar el checkpoint.
await client.query(
`update projection_checkpoint
set last_event_id = greatest(last_event_id, $1), updated_at = now()
where projection = 'orders_read'`,
[event.id],
)
await client.query(39;commit39;)
} catch (error) {
await client.query(39;rollback39;)
throw error
} finally {
client.release()
}
}
Lo importante: los tres pasos van en la misma transacción. Si el proceso muere entre el paso 1 y el 2, el rollback deshace la reclamación y el evento se vuelve a entregar. Sin transacción, esa tabla de deduplicación no te protege, te miente.
Este parseo de eventos es, por cierto, uno de los sitios donde Zod paga solo: el mismo schema te da el tipo de TypeScript, la validación en runtime y el mensaje de error que vas a leer en el dead letter a las tres de la mañana.
Si quieres exprimir esa parte —discriminated unions por event_type, transformaciones, versionado de schemas— lo trabajo a fondo en el curso de Zod para TypeScript.
Orden y versiones: cuando el v3 llega antes que el v2
Ningún bus te garantiza el orden global. Con particiones, reintentos y varios consumidores en paralelo, el evento v3 de un pedido puede llegar antes que el v2. Es normal, no es un bug del broker.
La cláusula que ya has visto arriba resuelve el caso:
where orders_read.version < excluded.version
Si llega el v3 y lo aplicas, cuando aparezca el v2 el WHERE da falso y el update no ocurre. El evento viejo se descarta en silencio, que es exactamente lo que quieres: tu read model no retrocede jamás.
Ahora la letra pequeña, que es donde se rompe la gente: esto solo funciona si tus eventos llevan el estado completo. Si order.paid dice "el total es 120 y el estado es paid", aplicar el v3 y tirar el v2 deja la vista correcta. Si tus eventos son deltas —"suma 3 al stock", "descuenta 20 del saldo"— descartar el v2 te deja con un número mal para siempre.
Con deltas necesitas detectar huecos y esperar. Cambias la guarda por una igualdad estricta:
-- solo aplico si soy exactamente el siguiente
where orders_read.version = excluded.version - 1
Y si rowCount === 0 y la versión del evento es mayor que la actual más uno, lanzas para que el bus te lo vuelva a entregar más tarde, cuando el que falta ya haya pasado.
Cambia también el SET: con deltas ya no copias excluded, acumulas — set total = orders_read.total + excluded.total. Y ojo al caso borde, que es el que muerde: esa guarda solo se evalúa en la rama DO UPDATE. Si la fila todavía no existe, el INSERT entra con la versión que traiga y el hueco pasa sin que nadie lo vea. Con deltas, crea la fila en la versión 0 cuando das de alta el agregado.
Mi recomendación después de sufrir las dos: haz los eventos state-carrying siempre que puedas. Pesan más en la cola y a cambio te ahorran toda la maquinaria de gaps, buffers y reentregas. Es el cambio de diseño más rentable de esta lista.
Rebuild de proyecciones: el superpoder que nadie usa
Aquí está la parte buena de CQRS, la que compensa todo lo anterior: si tu read model es una función pura de la secuencia de eventos, el read model es desechable. ¿Se corrompió por un bug del proyector? Lo tiras. ¿Quieres añadir una columna calculada a la vista? Lo tiras. ¿Cambias la forma entera de la proyección? Lo tiras.
Con la condición que casi nadie cumple: no borres los eventos. Un outbox con delete from outbox where published_at is not null es un outbox que funciona y que te quita esta capacidad para siempre. Archiva a un event_log en vez de borrar. Es la diferencia entre una cola y un log.
El patrón es proyección versionada, y son cuatro pasos:
- Creas
orders_read_v2con el esquema nuevo, vacía, y su propia fila enprojection_checkpoint. - Arrancas el proyector v2 en modo replay, leyendo el
event_logdesde el id 0. El v1 sigue vivo y sirviendo tráfico. - Cuando el v2 alcanza al v1 y ambos consumen en tiempo real, comparas. Unos cuantos agregados a mano o un diff de checksums.
- Cambias el puntero.
Ese cambio de puntero es lo único delicado. Si la capa de consulta lee a través de una vista, es una sentencia:
begin;
drop view orders_read_current;
create view orders_read_current as select * from orders_read_v2;
commit;
Dentro de la transacción toma un ACCESS EXCLUSIVE sobre la vista: las consultas en vuelo esperan unos milisegundos y siguen. Y no, create or replace view no vale aquí: solo admite añadir columnas al final, no cambiar nombres, tipos ni orden —que es exactamente lo que cambia en un rebuild.
Si no tienes vista, usa un flag de configuración que lea la capa de consulta al construir la query: más código, pero te deja volver atrás sin desplegar.
Con un rebuild fiable, tocar el read model deja de dar miedo. Ya no migras datos con un ALTER TABLE a las dos de la mañana: construyes una tabla nueva en paralelo, con tráfico real, y decides con datos si la enciendes.
Un apunte de método: el evento order.paid es una API pública aunque no tenga endpoint. Quién lo emite, qué campos garantiza y cómo se versiona tiene que estar escrito antes de picar el proyector.
Es de lo que más insisto en el libro de Spec-Driven Development, y en sistemas de eventos se nota el doble: el coste de equivocarte no lo pagas en el deploy, lo pagas seis meses después, cuando ya hay cuatro consumidores.
Medir el lag del read model: el único aviso temprano que vas a tener
En CQRS la consistencia eventual no es un fallo, es el contrato: el read model siempre va algo por detrás del write model. El fallo es no saber cuánto.
Todo lo anterior puede estar bien implementado y aun así tu read model puede ir veinte minutos por detrás porque el relay se quedó colgado. No se lanza ninguna excepción. No hay error 500. Todo está "verde".
Solo hay una métrica que te avisa: la antigüedad del evento pendiente más viejo.
select coalesce(
extract(epoch from now() - min(o.occurred_at)),
0
) as lag_seconds
from outbox o
left join processed_events p
on p.projection = 39;orders_read39;
and p.event_id = o.id
where p.event_id is null;
Devuelve 0 cuando no hay nada pendiente y crece cuando algo se atasca. Exponla como gauge y ponle alerta.
No la escribas contra last_event_id del checkpoint. Como el checkpoint avanza con greatest(), es una marca de agua alta: el v2 que se fue al dead letter queda por debajo de ella, la query no lo ve y te devuelve 0 con la proyección rota. Que es, literalmente, el ticket de los once días. Si purgas processed_events, limita el anti-join a tu ventana de retención.
Cuidado con la versión ingenua de esta métrica, que es la que suele estar puesta: now() - last_event_at del checkpoint. Esa te mide "cuánto hace que proyecté algo", y si a las tres de la madrugada no hay tráfico te va a despertar sin motivo. Peor: te acostumbra a ignorar la alarma. Mide lo que está esperando, no lo último que hiciste.
Yo añado dos series más al dashboard:
- Lag en eventos: cuántas filas del outbox siguen sin aparecer en
processed_events. Te dice si el proyector está perdiendo la carrera. No lo calcules comomax(id) − last_event_id: arrastra exactamente el mismo punto ciego de la marca de agua. - Tamaño del dead letter: si crece, hay eventos que no se están aplicando y el read model ya está mal.
Con esas tres, aquel ticket de los once días se habría abierto en once minutos.
Cuándo no necesitas absolutamente nada de esto
Si tu "read model" es una réplica de lectura de la misma base de datos, no tienes este problema. Postgres replica por ti, la sincronía la resuelve el WAL, y tu única métrica es el lag de replicación —que ya viene dado por pg_last_xact_replay_timestamp() y pg_stat_replication.
Nada de outbox, nada de proyectores, nada de checkpoints. Si separaste lectura y escritura solo para repartir carga, esa es la respuesta correcta y es aburridísima, que es justo lo que quieres en infraestructura.
Lo mismo si tu vista denormalizada es una materialized view en la misma base y toleras refrescarla cada pocos minutos. REFRESH MATERIALIZED VIEW CONCURRENTLY resuelve más casos de los que la gente cree. Pide dos cosas: un índice UNIQUE sobre columnas —sin expresiones y sin WHERE— y que la vista ya esté poblada. A cambio refresca sin bloquear lecturas, aunque tarda bastante más que un refresh normal.
Todo lo de este post empieza a hacer falta cuando el read model vive en otro sitio: otro motor, otro servicio, un índice de búsqueda, una tabla con una forma que no se deriva de un SELECT. Ahí sí tienes dual write y ahí sí necesitas el outbox.
| Dónde vive tu read model | Cómo se sincroniza | Qué tienes que operar | Lag típico |
|---|---|---|---|
| Réplica de lectura, misma base | Replicación física (WAL) | Nada, lo hace Postgres | Milisegundos |
| Materialized view, misma base | REFRESH MATERIALIZED VIEW CONCURRENTLY |
Un cron | Minutos |
| Otra tabla, otro servicio, índice de búsqueda | Transactional outbox + relay | Tabla outbox, relay, checkpoints |
Segundos |
| Igual que arriba, sin mantener relay | CDC leyendo el WAL | Kafka + Connect | Segundos |
Y si estás en ese caso pero no quieres mantener la tabla ni el relay, mira CDC antes de escribir código: resuelve lo mismo leyendo el WAL, a cambio de infraestructura extra. Lo desarrollo en las preguntas de abajo.
Qué hacer hoy
Si ya tienes CQRS en producción y nada de esto está montado, no empieces por el outbox. Empieza por la métrica.
Escribe la query del lag, ponla en un dashboard y déjala una semana. Vas a descubrir dos cosas: cuántos eventos estabas perdiendo sin saberlo, y si tu problema real era ese o era otro. Es media hora de trabajo, y es lo único de esta lista que te da información antes de que te la pida un cliente enfadado. El resto —outbox, idempotencia, versiones, rebuild— se construye después, con datos encima de la mesa.
Si quieres ver este tipo de arquitecturas montadas de principio a fin, con el código completo y las decisiones discutidas, en Dominicode Labs es donde publico los proyectos largos que no caben en un post.
Preguntas frecuentes
¿Por qué mi read model no se actualiza en CQRS?
Casi siempre por dual write: el cambio de estado se guardó, pero el evento nunca llegó al bus porque publicar y hacer commit son dos operaciones distintas y la segunda falló sin que nadie reintentara.
Los otros dos sospechosos habituales son el relay parado —el evento sigue en la tabla outbox con published_at a null— y un evento que el proyector rechaza una y otra vez hasta acabar en el dead letter. Los tres casos se distinguen en treinta segundos con la query de lag de este post: si devuelve un número alto, el evento existe y no se ha proyectado; si devuelve 0 y la vista sigue mal, mira el dead letter.
¿El transactional outbox añade latencia a cada escritura?
Añade un INSERT dentro de una transacción que ya estaba abierta. En la práctica es ruido comparado con el resto del comando.
La latencia que sí importa es la otra: cuánto tarda el evento en llegar al read model. Eso lo marca el intervalo de polling del relay, no el insert. Si necesitas bajarlo, usa LISTEN/NOTIFY de Postgres para despertar al relay en cuanto hay una fila nueva, en lugar de esperar al siguiente tick.
¿Puedo usar CDC en lugar de la tabla outbox?
Sí, y resuelve el mismo problema de dual write. Debezium lee el WAL de Postgres y publica los cambios sin que tu código haga nada — trae incluso un outbox event router preparado exactamente para este patrón.
La diferencia es qué publicas. Con outbox publicas eventos de dominio que tú diseñas; con CDC a secas publicas cambios de filas, y tus consumidores acaban acoplados al esquema de tu base de datos. El punto medio más usado es CDC leyendo precisamente la tabla outbox: eventos de dominio sin escribir relay, a cambio de operar Kafka y Connect.
¿Qué hago con un evento que el proyector nunca consigue procesar?
Dead letter después de N intentos, y alerta. Lo que no puedes hacer es reintentarlo en bucle para siempre: bloqueas la partición y frenas todo lo que viene detrás.
Ojo con la consecuencia que se pasa por alto: si mandas a dead letter el v2 de un agregado y sigues procesando el v3, esa fila queda incoherente hasta que reproceses. Por eso el dead letter va en el dashboard y no en un buzón que nadie abre.
¿Necesito event sourcing para poder reconstruir proyecciones?
No. Necesitas retener los eventos, que es mucho menos que event sourcing.
En event sourcing el log de eventos es la fuente de verdad y el estado se deriva de él. Aquí la fuente de verdad sigue siendo tu tabla orders, y el log de eventos es solo el historial de cambios publicados. Con archivar el outbox en un event_log en vez de borrarlo ya puedes reconstruir cualquier proyección.
¿Cada cuánto debería hacer polling del outbox?
Depende del lag que tu producto tolere, no de lo que haga la industria. Un panel interno aguanta segundos; un contador que el usuario ve moverse tras pulsar un botón, no.
Define el número primero —"el read model va como mucho X segundos por detrás"—, mídelo con la query de lag y ajusta el intervalo hasta cumplirlo. Sin ese número escrito, cualquier valor que pongas es una opinión.
Por Bezael Pérez — Developer senior con más de 15 años de experiencia y fundador de Dominicode.
¿Te resultó útil este artículo?
Compártelo con tu comunidad y ayuda a otros desarrolladores.
