Published

🎯 2d — Flip global de pipeline.output a Iceberg nativo (approach de espectro completo)

Connect any source, model it as an ontology, transform it, and operationalize it, analytics, automation and machine learning, under one governed, self-hostable roof. --- Most teams stitch the...

🎯 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_rows y 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 previoRealidad 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_rowsun 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

ProductorFicherowriterId / seamEstrategiasEstado
Deploy materialiser (el pipeline.output canónico)lib/workers/pipelines/dataset-output-materialiser.ts (:319 seam) + deploy-worker.tsresolveDatasetRowSink('pipeline.output')maxMode=pg, sin brazo icebergdataset-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:415NINGUNOINSERT INTO dataset_rows crudo, no toca el seaminline (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

  1. Romper el WITH…INSERT — separar compute (PG, produce un AsyncIterable<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.
  2. 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).
  3. 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).
  4. El delta nativo + serialización — el deltaPredicate sobre dataset_rows no 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.expectations source-respecting. ✅ HECHO (2026-07-04). Las 4 checks (rowCount, primaryKey, rowLevel, value/floats vía _shared) resuelven la fuente por resolveSqlSource + hydratedFrom (el mismo hydrate-to-inline-jsonb de sql-editor/dashboards-QUERY, en _shared.resolveExpectationSource). Hoy PG vía ctx.query (read-your-writes en la tx del materialise); un output nativo hidrata la snapshot a jsonb_to_recordset y la agregación (COUNT FILTER / GROUP BY / predicados compilados) corre idéntica → ya NO valida en falso sobre un dataset_rows drenado.
  • P0.2 — Rutar el source read de model-node por el linchpin. ✅ HECHO (2026-07-04). model-node-materialiser.ts ya resuelve su upstream por resolvePipelineSource (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.ts transform/join/split + runtime.ts) rutados por resolveDatasetRowSink('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 de ensureProducedDataset, 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.ts gana 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 open sin committed tras un timeout = build abortado a mitad → un barredor la marca aborted y 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 / nodoPor 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 deploysin filas previas → degenera a replace puro.
append_alwaysappend; 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):

EstrategiaTarget-read a rutar por facadeComplicació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) + DELETEEl 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 (checkflipbackfillverify por checksum PG-vs-Iceberg). CUTOVER_WRITERS NO incluye pipeline.output (:44) → extenderlo.
  • lib/lakehouse/canary.ts + lakehouse_read_parity: paridad PG vs Iceberg por dataset (counts + spot-check por __row_index). Solo un MISMATCH caught-up bloquea. TTL de flag = 10 s.

Secuencia canary recomendada (por-dataset):

  1. Elegir hojas del DAG (sin downstream) con estrategia replace-pura y sin expectations hasta cerrar P0.1.
  2. 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).
  3. lakehouse-canary.ts --smoke --dataset <ds>: exige ok/no-MISMATCH.
  4. Flipear los readers del output al final (lakehouse-flags.ts set <reader> iceberg --dataset <ds>) — todos ya tienen ejecutor iceberg salvo pipeline.expectations (P0.1).
  5. Revertir = clear/set pg --dataset <ds> (≤10 s). NO drenar el dataset_rows del 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 maxMode o el motor. El flip real (2d.1+) es una tanda posterior, sobre terreno seguro.

Abiertas (para cuando lleguemos ahí):

DecisiónOpcionesRecomendación
Compute→streamcursor PG (FETCH) vs materializar a temp vs SELECT completo a memoriaCursor PG (memory-bounded, como el stream de la facade); evita drenar a memoria y no necesita temp.
Merge-preservanteleer target por facade + merge en Node vs primitiva merge en ml-runnerLeer por facade + merge en el productor (reusa la lógica SQL existente traducida); una primitiva ml-runner es optimización posterior.
schema-reconcile nativoUPDATE in-place (imposible en Iceberg) vs poda en el compute vs replace lo obviaEn 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 retryreplace idempotente vs 2-phase-commit realReplace 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 paraleloP0.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.output sin 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_index entre builds concurrentes.
  • ❌ Delta nativo por deltaPredicate sobre dataset_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=buildStart de reloj local → skew vs committed_at nativo → delta saltado permanentemente.
  • ❌ Drenar dataset_rows del 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.