🎯 2d — Flip global de pipeline.output a Iceberg nativo (approach de espectro completo)
Objetivo. Que el WRITE de los outputs de pipeline deje de aterrizar en Postgres
dataset_rowsy escriba directamente al tridente (R2 + Lakekeeper + PG-metadata) vía
facade.write()/ ml-runner. Es el gran proyecto acoplado del paso 6;
el linchpin READ ([source-resolver.ts]) y todo el 2c (readers analytics) ya están cerrados,
así que el camino de LECTURA está plano — pero el WRITE tiene su propio espectro.Fundado en 5 auditorías sobre el terreno (2026-07-04). Este doc CORRIGE varias
suposiciones que arrastraba paso6-write-approach.md — ver §0.
Companion: port-migration.md (tracker por-puerto) · dataspace-router.md.
0. Correcciones auditadas (los mitos que la auditoría desmintió)
El approach previo daba por hechas cosas que el código NO respalda. Antes de diseñar, corregir el mapa:
| # | Mito en el doc previo | Realidad auditada |
|---|---|---|
| M1 | "la escritura es una tx atómica del build" | Atomicidad es por-output, no por-deploy (deploy-worker.ts:582-681, BEGIN/COMMIT por output, sin savepoint — el savepoint solo existe en el Apply de diseño de transformPaths.ts:243). Si el output #3 falla, #1/#2 quedan commiteados (partial-build tipo Foundry). |
| M2 | "append-desde-max es seguro porque los builds por dataset están serializados — ya lo están" | FALSO. No hay advisory-lock ni FOR UPDATE sobre el dataset producido en NINGÚN camino de build (grep: 0 matches). Deploy corre concurrency=3 (deploy-worker.ts:1157), dedup solo por deploymentId. Hoy el DELETE+INSERT atómico de PG enmascara la ausencia de lock (dos replaces concurrentes = last-wins, sin corrupción). La serialización es un prerrequisito a CONSTRUIR, no un hecho. |
| M3 | "el delta ledger es storage-agnostic (solo metadata txn)" | MIXTO. La resolución de la ventana (low, high] sí es agnóstica (lee committed_at/type vía la view dataset_transactions_all). Pero su aplicación (deltaPredicate, incremental.ts:70-83) genera transaction_id IN (…) sobre un scan FROM dataset_rows → un source nativo con dataset_rows drenado no tiene ese transaction_id por-fila → el delta-append lee CERO filas en silencio. |
| M4 | "camino plano — todos los readers son source-respecting" | Verdadero para readers interactivos/analytics. Dos excepciones que SÍ leen un output de pipeline: (a) pipeline.expectations es NO-EXEC (seam-wired sin ejecutor iceberg) → sobre un output nativo lee dataset_rows drenado y valida en falso (rowCount=0 pasa); (b) model-node-materialiser.ts:199 lee su source upstream con FROM dataset_rows crudo, sin resolvePipelineSource. |
| M5 | (implícito) "el ratchet garantiza que cada acceso está rutado" | El ratchet es file-level: un fichero con CUALQUIER resolveDatasetRowS(ink|ource) queda exento entero. Oculta accesos crudos en ficheros seam-aware (así se esconde model-node-materialiser.ts:199). |
| M6 | "facade.write({upsert}) cubre los merges" | El upsert de PyIceberg (writer.py:388, merge-on-read por PK: update matched + insert new) NO es la semántica de snapshot_replace (que preserva filas del target cuyo PK no está en el batch, con last-wins intra-batch) ni de append_new (anti-join). Gap real. |
| M7 | (implícito) "snapshot_replace es replace" | Solo lo es en el primer deploy (sin filas previas). Con filas previas lee y preserva el target → es merge-preservante (snapshot-replace.ts:120-123). |
Además, dos primitivas de build no tienen equivalente nativo: schema-reconcile.ts (UPDATE dataset_rows in-place que poda claves JSONB stale — Iceberg no hace UPDATE in-place) y checkpoint-ephemeral.ts (clona el estructural PG→PG con INSERT…SELECT FROM dataset_rows → 0 filas si el estructural es nativo → aislamiento de build roto). Y el estado open de pipeline_build_transactions existe en el CHECK pero está muerto (el pipeline siempre inserta committed inline) → el write-intent ledger lo revivirá.
1. El terreno: TRES productores, DOS mundos de escritura
| Productor | Fichero | writerId / seam | Estrategias | Estado |
|---|---|---|---|---|
Deploy materialiser (el pipeline.output canónico) | lib/workers/pipelines/dataset-output-materialiser.ts (:319 seam) + deploy-worker.ts | resolveDatasetRowSink('pipeline.output') — maxMode=pg, sin brazo iceberg | dataset-write-strategies/* (default/snapshot_replace/append_new/append_always) | 🔴 clamp PG |
| Transform-paths leak (Apply interactivo del canvas) | lib/pipelines/transformPaths.ts (transform :341, join :747, split :1295/:1303) + legacy runtime.ts:415 | NINGUNO — INSERT INTO dataset_rows crudo, no toca el seam | inline (DELETE+re-insert / append-delta) | ⛔ bypass total |
| Model-node materialiser (predictions) | lib/workers/pipelines/model-node-materialiser.ts (:241 seam, :199 source raw) | seam solo-PG + source read crudo (:199, no por linchpin) | INSERT jsonb por ventana; row_index reusa el del source | 🔴 clamp PG + source raw |
Los tres aterrizan en dataset_rows. El object-type output (object-type-output-materialiser.ts → tabla objects) NO es parte de este flip (su substrato es objects, no un dataset Iceberg).
El write real (todos los productores): un único WITH … [DELETE] … INSERT INTO dataset_rows donde compute (read literal FROM dataset_rows) y write comparten statement y tx. Solo el sink pasa por el seam (y clampa a PG). El row_index lo asigna el productor: row_number() (replace) o COALESCE(MAX(row_index),-1)+row_number() (append, lee el target).
2. Los cuatro sub-problemas del flip
- Romper el
WITH…INSERT— separar compute (PG, produce unAsyncIterable<Row>de filas de negocio SIN identidad) del write (facade.write({kind:'rows', rows, columns, writeMode, sourceRefs})→ ml-runner estampa__row_index/__row_id/__created_at). El compute deja de INSERTar; materializa a un cursor/stream. - Recuperar la atomicidad — compute PG (rollbackable) + write nativo (commit externo a R2/Lakekeeper) dejan de ser una tx. Se resuelve con un write-intent ledger (§4).
- Los target-reads por facade — los modos merge-preservantes leen el target existente (
FROM dataset_rows WHERE dataset_id=$2); en nativo está drenado → esa lectura va por la facade (§5). - El delta nativo + serialización — el
deltaPredicatesobredataset_rowsno funciona en nativo, y el append-desde-max exige serialización que hoy no existe (§6).
3. Prerequisitos BLOQUEANTES (P0 — antes de subir el techo maxMode)
Ninguno de estos es "el flip"; son los huecos que la auditoría destapó y que harían fallar en silencio un output nativo. Cerrarlos ANTES de dar un ejecutor iceberg a pipeline.output:
- P0.1 —
pipeline.expectationssource-respecting. ✅ HECHO (2026-07-04). Las 4 checks (rowCount, primaryKey, rowLevel, value/floats vía_shared) resuelven la fuente porresolveSqlSource+hydratedFrom(el mismo hydrate-to-inline-jsonb de sql-editor/dashboards-QUERY, en_shared.resolveExpectationSource). Hoy PG víactx.query(read-your-writes en la tx del materialise); un output nativo hidrata la snapshot ajsonb_to_recordsety la agregación (COUNT FILTER / GROUP BY / predicados compilados) corre idéntica → ya NO valida en falso sobre undataset_rowsdrenado. - P0.2 — Rutar el source read de model-node por el linchpin. ✅ HECHO (2026-07-04).
model-node-materialiser.tsya resuelve su upstream porresolvePipelineSource(FROM ${sourceTable}) como transform/join/split/runtime → PG no-op hoy, hidrata nativo mañana. - P0.3 — Cerrar la fuga de transform-paths por el seam. ✅ HECHO (2026-07-04). Los 4 INSERT (
transformPaths.tstransform/join/split +runtime.ts) rutados porresolveDatasetRowSink('pipeline.output').run({postgres})→ seam-aware + flippable (PG clamp hoy). El split (dual-target) lanza si una rama se flipea antes de 2d.1. - P0.4 — Serialización por dataset producido. ✅ HECHO (2026-07-04).
lib/pipelines/dataset-build-lock.ts(lockDatasetForWrite/lockDatasetsForWrite) =pg_advisory_xact_lock(ns, hashtext(id))xact-scoped, cableado en los 5 sitios de write: deploy dataset-output (antes del delta-read + strategy), model-node (dentro deensureProducedDataset, antes del DELETE-on-reuse + tras el mint), y las 3 applies (transform/join = 1 lock; split = 2 en orden ordenado). Deadlock-free por construcción (≤1 lock por tx; split ordena). Runtime legacy exento (fresh-dataset-per-run, uncontended). - P0.5 — Ratchet a access-level. ✅ HECHO (2026-07-04).
check-dataset-rows-seam.tsgana un count-ratchet: pinnea el nº de accesos crudos de CADA fichero exento (allowlist ∪ seam-aware, 43 ficheros) y falla si SUBE (un nuevo bypass colado en un fichero ya exento — el blindspot file-level de M5). Regenerar el mapa:SEAM_BASELINE_REPORT=1 npm run check:dataset-rows-seam.
4. El write-intent ledger (recuperar atomicidad)
Hoy pipeline_build_transactions se inserta committed inline en la misma tx que el INSERT de filas (dataset-output-materialiser.ts:342, transformPaths.ts:322). Al externalizar el write, revivir el ciclo open→commit que ya usa el motor nativo (iceberg-native-write.ts:127-157) y que el CHECK de la tabla ya admite (estado open, hoy muerto):
1. INTENT — INSERT pipeline_build_transactions (status='open', high=buildStart) [tx PG corta]
2. COMPUTE — compute PG → AsyncIterable<Row> de filas de negocio [rollbackable]
3. WRITE — facade.write({rows, writeMode:'replace', sourceRefs}) [ml-runner, commit externo]
4. COMMIT — UPDATE pipeline_build_transactions status='committed', row_count [tx PG corta]
+ UPDATE datasets.source_watermark = buildStart (solo AQUÍ)
+ UPDATE datasets.row_count/status
Invariantes que preserva:
- El watermark solo avanza tras write confirmado (paso 4). Avanzarlo antes perdería el delta si el write falla.
- Idempotencia por replace — si el write commitea pero el paso 4 falla, el retry re-ejecuta el
write({replace})(ml-runner posee row_index, reescribe todo) → seguro. El append-desde-max NO es idempotente tras un commit parcial → por eso replace-first (§5). - Recovery: una fila
opensincommittedtras un timeout = build abortado a mitad → un barredor la marcaabortedy el dataset queda en su último estado committeado (el write nativo replace es all-or-nothing a nivel snapshot Iceberg).
⚠️ Ojo al clock-skew de la ventana half-open (M3-relacionado): high=buildStart es wall-clock del worker; el committed_at del write nativo lo estampa otro proceso. Si el nativo commitea con committed_at > buildStart del downstream, esa tx cae fuera de (low, high] y se salta permanentemente. Fix: derivar high del committed_at REAL devuelto por facade.write() (el CommitResult trae transactionId), no del reloj local.
5. Replace-first: la secuencia por estrategia
facade.write({replace}) mapea limpio cuando el productor REGENERA el output entero y ml-runner posee el row_index. Los merge-preservantes necesitan leer el target por la facade primero.
Tier A — replace-first (fáciles, primero; NO leen el target):
| Estrategia / nodo | Por qué es limpio |
|---|---|
join (applyJoinPath) | DELETE + re-insert full, sin target-read. El más limpio. |
split (applySplitPath + compilers/split.ts) | DELETE ambas ramas + re-insert particionado (partición exhaustiva). |
transform snapshot (!isDelta) | DELETE + re-insert full. |
| model-node (predictions) | DELETE-on-reuse + re-insert por ventana (ojo: preservar la correlación SOURCE_ROW_KEY fila-source↔predicción). |
| snapshot_replace / default — PRIMER deploy | sin filas previas → degenera a replace puro. |
| append_always | append; ml-runner continúa __row_index desde el total de la tabla (writer.py:334). Solo suelta el MAX(row_index) read — disuelto por ml-runner. |
Tier B — requieren target-read por facade (después):
| Estrategia | Target-read a rutar por facade | Complicación |
|---|---|---|
| append_new | __existing_pks__ FROM dataset_rows WHERE dataset_id=$2 (append-new.ts:77-80) | Leer los PKs del target por la facade para el anti-join, luego write({append}) con las filas nuevas. |
| snapshot_replace / default con filas previas | __preserved_candidates__ FROM dataset_rows WHERE dataset_id=$2 (snapshot-replace.ts:120-123) + DELETE | El caso más duro. Leer TODAS las filas preservadas del target por la facade, mergear con el batch nuevo en el productor (last-wins por PK), y write({replace}) el conjunto combinado. El upsert de PyIceberg NO sirve (M6). |
| transform-delta (incremental) | offset MAX(row_index) (transformPaths.ts:339) + el deltaPredicate (§6) | Doble problema: el delta window + el offset, ambos sobre dataset_rows. |
Row_index: todos los row_number() puros (replace) desaparecen — los posee ml-runner. Los offsets MAX(row_index) del target (append) los gestiona ml-runner del lado nativo (continúa desde el total) siempre que los builds estén serializados (P0.4).
6. El sub-problema más sutil: el delta nativo
El deltaPredicate (incremental.ts:70-83) filtra transaction_id IN (SELECT … FROM dataset_transactions_all WHERE committed_at ∈ (low,high]) sobre un scan FROM dataset_rows. En un source nativo, dataset_rows está drenado y no hay transaction_id por-fila → el delta lee 0 filas en silencio (M3, y el anti-patrón "mixed-mode delta" de paso6:117).
Opciones (decisión abierta — ver §9):
- (a) Delta por metadata txn nativa — el snapshot Iceberg lleva
__created_at/txn; resolver la ventana contra la metadata del ledger nativo (dataset_transactions) y leer el delta por la facade con un filtro temporal. Requiere que ml-runner soporte un read filtrado por txn/tiempo. - (b) Clamp delta a PG mientras el output es nativo — pero eso ES mixed-mode (el riesgo que queremos evitar); solo válido si el downstream que lee ese output en delta NO se ha flipeado. Frágil.
- (c) Full-rebuild-first — deshabilitar el modo incremental para un output nativo hasta tener (a); el output nativo siempre se regenera entero (replace). Simple, correcto, pierde la optimización incremental. Recomendado para el arranque (replace-first ya implica esto).
Serialización (P0.4) es prerrequisito del append/delta nativo: sin lock, dos builds concurrentes con write externo no-atómico colisionan el row_index.
7. Flags, maxMode, canary y observabilidad
Subir el techo: writer-flags.ts:109-112 pipeline.output defaultMode:'pg'→… y maxMode:'pg'→'iceberg_native' (hoy el clamp lo pinnea a PG incluso con flag; el CLI lo rechaza). Añadir el brazo iceberg: al sink.run() del materialiser + estrategias (patrón item-consumption.ts:601-625). Mantener defaultMode:'pg' al principio → flip opt-in por-dataset (canary), como fue manual.table en 2a (que sí subió ambos).
Infra de cutover por-dataset (existe, sólida, pero calibrada para Tier-1):
scripts/lakehouse-write-flags.ts(pg|dual|iceberg_native,--dataset),scripts/lakehouse-flags.ts(readers),scripts/lakehouse-cutover.ts(check→flip→backfill→verifypor checksum PG-vs-Iceberg).CUTOVER_WRITERSNO incluyepipeline.output(:44) → extenderlo.lib/lakehouse/canary.ts+lakehouse_read_parity: paridad PG vs Iceberg por dataset (counts + spot-check por__row_index). Solo unMISMATCHcaught-up bloquea. TTL de flag = 10 s.
Secuencia canary recomendada (por-dataset):
- Elegir hojas del DAG (sin downstream) con estrategia replace-pura y sin expectations hasta cerrar P0.1.
lakehouse-cutover.ts check <ds>(checksum PG vs mirror) →set pipeline.output iceberg_native --dataset <ds>→ re-deploy (escribe snapshot nativo) →verify <ds>(checksum PG vs snapshot).lakehouse-canary.ts --smoke --dataset <ds>: exigeok/no-MISMATCH.- Flipear los readers del output al final (
lakehouse-flags.ts set <reader> iceberg --dataset <ds>) — todos ya tienen ejecutor iceberg salvopipeline.expectations(P0.1). - Revertir =
clear/set pg --dataset <ds>(≤10 s). NO drenar eldataset_rowsdel output hasta que el canary sostenga paridad varios ciclos (revertir tras drenar exige backfill inverso — eso es Fase 5).
Gaps de observabilidad a cerrar: (a) el fallback silencioso NO-EXEC del read-seam no se loguea (read-router.ts:206-213 cae a postgres() sin traza) → un output nativo leído por un NO-EXEC no deja rastro; añadir un warn. (b) No hay parity table write-side análoga a lakehouse_read_parity (la evidencia write es el checksum on-demand del cutover, no continua) → considerar una.
8. Secuenciación por fases
P0 · PRERREQUISITOS ✅ COMPLETO (2026-07-04) — source-respecting + serialización
├─ P0.1 pipeline.expectations → hydrate-to-inline-jsonb ✅
├─ P0.2 model-node source read → linchpin ✅
├─ P0.3 transform-paths WRITE → resolveDatasetRowSink ✅
├─ P0.4 advisory-lock por dataset producido (deploy + Apply) ✅
└─ P0.5 ratchet access-level (count-ratchet) ✅
⇒ terreno listo para 2d.1 (el motor). Ya NINGÚN reader/writer de pipeline
lee/valida un dataset_rows drenado; los builds por dataset están serializados.
2d.1 · MOTOR DEL PRODUCTOR (romper el WITH…INSERT) → APPROACH DETALLADO (5 auditorías):
│ docs/architecture/2d1-engine-approach.md
├─ compute PG → AsyncIterable<Row> (cursor pg-query-stream, sin INSERT)
├─ brazo iceberg: en el sink del materialiser + estrategias → facade.write({rows})
├─ write-intent ledger (revivir open→commit; serializar por unique-index status='open',
│ NO el advisory-lock xact-scoped que no sobrevive la partición de tx; committedAt real
│ a construir — el CommitResult NO lo trae hoy)
├─ FIX bloqueantes: recount transform-path → CommitResult · cross-join count → handle.rowCount
└─ subir maxMode pipeline.output → iceberg_native (defaultMode sigue pg)
2d.2 · REPLACE-FIRST (Tier A, canary por-dataset)
└─ join · split · transform-snapshot · model-node · append_always · snapshot_replace-1er-deploy
(delta: modo (c) full-rebuild-first — sin incremental nativo aún)
2d.3 · MERGE-PRESERVANTE (Tier B, target-read por facade)
├─ append_new (leer PKs del target por facade → anti-join → write append)
├─ snapshot_replace preservante (leer filas preservadas por facade → merge → write replace)
└─ delta nativo real (opción (a): ventana por metadata txn nativa + read filtrado)
DIFERIDO / FASE 5
├─ drenar dataset_rows de outputs nativos (tras N ciclos de paridad)
└─ schema-reconcile / checkpoint-ephemeral nativos (sin equivalente hoy)
9. Decisiones
TOMADAS (2026-07-04):
- ✅ Delta nativo = (c) full-rebuild-first. Un output nativo se regenera entero (replace); el incremental nativo se difiere a 2d.3 (opción (a), metadata txn nativa). Encaja con replace-first. Nunca (b) mixed-mode.
- ✅ Arranque = SOLO P0 primero. Implementar + verificar los 5 prerrequisitos bloqueantes (§3) ANTES de tocar
maxModeo el motor. El flip real (2d.1+) es una tanda posterior, sobre terreno seguro.
Abiertas (para cuando lleguemos ahí):
| Decisión | Opciones | Recomendación |
|---|---|---|
| Compute→stream | cursor PG (FETCH) vs materializar a temp vs SELECT completo a memoria | Cursor PG (memory-bounded, como el stream de la facade); evita drenar a memoria y no necesita temp. |
| Merge-preservante | leer target por facade + merge en Node vs primitiva merge en ml-runner | Leer por facade + merge en el productor (reusa la lógica SQL existente traducida); una primitiva ml-runner es optimización posterior. |
schema-reconcile nativo | UPDATE in-place (imposible en Iceberg) vs poda en el compute vs replace lo obvia | En replace el re-write limpia solo (obvia el prune). En append-preservante hay que podar las claves stale en el merge del productor. |
| Atomicidad del retry | replace idempotente vs 2-phase-commit real | Replace idempotente (write-intent ledger). 2PC real es sobre-ingeniería — ml-runner ya es all-or-nothing por snapshot. |
| ¿P0 antes o durante? | cerrar P0 completo antes de 2d.1 vs en paralelo | P0.1/P0.2/P0.4 ANTES (rompen en silencio). P0.3/P0.5 pueden ir con 2d.1. |
10. Anti-patrones (que este approach evita)
- ❌ Flipear
pipeline.outputsin cerrar P0.1 → expectations validan en falso sobre 0 filas. - ❌ Flipear sin cerrar P0.2 → un modelo lee su upstream nativo drenado (0 filas) y "entrena/predice" vacío.
- ❌ Append-desde-max sin serialización (P0.4) → colisión de
row_indexentre builds concurrentes. - ❌ Delta nativo por
deltaPredicatesobredataset_rows→ lee 0 filas en silencio (M3). - ❌ Avanzar el watermark antes de confirmar el write nativo → pérdida de delta si el write falla.
- ❌
high=buildStartde reloj local → skew vscommitted_atnativo → delta saltado permanentemente. - ❌ Drenar
dataset_rowsdel output antes de sostener paridad → revertir exige backfill inverso. - ❌ Confiar en el ratchet file-level como prueba de cobertura (M5).
Approach de espectro completo, fundado en 5 auditorías (2026-07-04). Mantener al día conforme cada estrategia aterrice en el tridente. El write sigue al linchpin READ (cerrado) y al 2c (cerrado); esto es lo que queda del paso 6.