Published

🔺 2d.3 · B4 — Delta nativo (writes incrementales al target Iceberg)

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.3 · B4 — Delta nativo (writes incrementales al target Iceberg)

Entregable de approach (2026-07-09). El caso B4 de 2d3-tier-b-approach.md: un pipeline incremental (readMode=incremental) cuyo OUTPUT es nativo. Hoy el brazo iceberg del deploy lanza native delta write not supported (full-rebuild-first) → todo output incremental cae a PG. Grounded en lib/pipelines/incremental.ts, source-resolver.ts, las estrategias, y el dispatch B1/B2/B3 ya cerrado.

Tesis del doc: B4 es más pequeño de lo que parece en código pero merece el espectro completo porque el delta tiene una asimetría real — la VENTANA es un concepto del ledger PG, la APLICACIÓN ya es nativa. Separar los dos ejes es todo el diseño.

✅ IMPLEMENTADO (opción A, 2026-07-09) — canary §5d verde E2E. computeReplaceSql delta-aware en las 3 estrategias (snapshot-replace/append-always/append-new: AND ${deltaPredicate('$1','$2','$3')} cuando ctx.sourceDelta) + throw del brazo iceberg retirado. El dispatch de re-deploy aplica la ventana (upsert/append/append_new). Verificado: source en 2 txns, la ventana (srcTxn1, srcTxn2] deja entrar SOLO la 2ª txn al output nativo. Asunción dataset_rows-no-drenado documentada como TODO(Fase 5) (§5), guard NO añadido (sería conservador en la ventana canary).


1. Qué es el delta hoy (recap grounded)

  • Ventana (low, high] — low = datasets.source_watermark del output (o EPOCH en el primer build incremental), high = buildStart. Acotar por high=buildStart hace el read race-safe (una txn que commitea durante el build la lee el SIGUIENTE, nunca a medias).
  • deltaPredicate(dsParam, lowParam, highParam) (incremental.ts:70) = narra el scan del SOURCE a las filas nuevas:
    transaction_id IN (
      SELECT transaction_id FROM dataset_transactions_all
       WHERE dataset_id = $1 AND committed_at > $low AND committed_at <= $high
    )
    
    Se AND-ea sobre WHERE dataset_id = <source> — solo narra el SOURCE, nunca el read del previous-output.
  • windowAllowsDelta (incremental.ts:109) — la ventana es delta-appendable SĂ“LO si TODA txn es APPEND/UPDATE. Un SNAPSHOT en la ventana ⇒ el source fue REESCRITO ⇒ el build cae a full rebuild (sourceDelta=undefined). Esta es la garantĂ­a de correctitud clave: el delta es puramente ADITIVO.
  • El linchpin FUERZA PG para el delta (source-resolver.ts:74, if (… || opts?.hasDelta) return PG_SOURCE): la ventana vive en el ledger PG (transaction_id + dataset_transactions_all.committed_at), sin equivalente en un snapshot Iceberg → el read del delta se queda en dataset_rows.

2. Los DOS ejes ortogonales (el corazĂłn de B4)

EjeQué esEstado
READ del deltaleer las filas del source committeadas en (low,high]🔒 hoy sobre PG (dataset_rows + dataset_transactions_all) — difícil de nativizar (el ledger es PG)
WRITE del deltaaplicar esas filas al target Icebergâś… YA hecho (B1/B2/B3 = upsert/append/append_new nativo)

Insight central: el native delta es el MISMO dispatch de re-deploy (nativeReplayMode/nativeAppendNew) con el compute leyendo la ventana en vez del scan completo. El throw está ANTES del dispatch. El delta NO cambia el eje WRITE — solo cambia QUÉ filas entran.

Mapeo readMode(incremental) Ă— writeMode:

  • incremental + snapshot_replace → las filas de la ventana mergean por PK en el output = upsert.
  • incremental + append_always → appendan = append.
  • incremental + append_new → appendan las de PK nuevo (anti-join) = nativeAppendNew.

Los tres write-modes YA existen (B1/B2/B3). El primer build incremental (low=EPOCH, target vacío) → el probe#0 no ve snapshot → replace de todo hasta high (= el "full append into empty" de planArtifactBuild). ✅


3. El espectro del READ (la Ăşnica parte no trivial)

OpciónCómo lee la ventanaCosteCuándo
A · Interim (recomendada)PG (dataset_rows, como hoy) — el compute añade el deltaPredicatemínimo (reusa todo)AHORA — mientras el source tenga dataset_rows legible
B · SOTA native-scanIceberg incremental scan snapshot→snapshot (mapear watermark↔snapshot_id)alto (nuevo read-path + tabla de mapeo watermark/snapshot)cuando se drene dataset_rows (Fase 5)
C · Karmael motor subsume delta + cursor nativo—era Karma

Por qué A es suficiente ahora: un output nativo con delta lee la ventana de su SOURCE, y el source hoy conserva sus dataset_rows (nada se ha drenado — misma disciplina "no drenar hasta paridad" sobre la que descansa TODO el canary 2d.x). El caso que rompe A — un source nativo Y drenado — no existe hasta Fase 5, para cuando B/C deberían haber aterrizado.


4. La opción A — cambios de código (ordenados, mínimos)

  1. computeReplaceSql delta-aware en las 3 estrategias (snapshot-replace.ts, append-always.ts, append-new.ts): añadir AND ${deltaPredicate('$1','$2','$3')} al WHERE dataset_id=$1 cuando ctx.sourceDelta está presente, y bindear low/high a $2/$3 — mirror exacto de lo que ya hace apply (que usa $4/$5 porque lleva el prefijo [source, target, txn]; el compute solo lleva [source], así que aquí son $2/$3). Un output nativo comparte la misma dedup/proyección; solo se narra la ventana.
  2. Quitar el if (sourceDelta) throw del brazo iceberg (dataset-output-materialiser.ts). El dispatch (probe#0 + nativeReplayMode/nativeAppendNew) ya hace lo correcto.
  3. NADA más en el WRITE — writeOutputToTrident + los write-modes nativos ya están. El outputTxnType='APPEND' del ledger para incremental ya se estampa.

Correctitud: windowAllowsDelta garantiza que ningún SNAPSHOT entra en la ventana → el delta es aditivo → upsert/append es seguro (no hay que reconciliar un source-replace). El re-deploy vs primer-build lo decide el probe#0. La atomicidad la da withNativeControlPlane (snapshot all-or-nothing) como en B1/B2/B3.


5. La ASUNCIÓN documentada (y el guard que NO añadimos aún)

El read de la ventana sigue en dataset_rows (PG). Dependencia explĂ­cita: funciona mientras dataset_rows del SOURCE no se drene.

  • Hoy (canary window): nada drenado → un source nativo aĂşn tiene su mirror PG → el deltaPredicate lee bien. âś…
  • Fase 5 (drain): un source nativo drenado → dataset_rows vacĂ­o → el deltaPredicate lee 0 filas en silencio (delta vacĂ­o = no-op incorrecto). ❌ AhĂ­ entra la opciĂłn B/C.

Guard opcional (recomendación: DIFERIR): detectar "source nativo con dataset_rows drenado" y fail-loud. NO lo añadimos aún porque en la ventana canary (mirror presente) sería demasiado conservador — bloquearía deltas que hoy funcionan. Se documenta como TODO(Fase 5): al drenar, native delta exige el read nativo de la ventana (opción B) o Karma (C). El guard se activa junto con el drain, no antes.


6. Verificación (canary — extender §5)

scripts/dataspaces/pipeline-output-canary.ts gana §5d (delta), mismo patrón (setup en main.test → write → read-back → paridad, limpieza):

  • Sembrar el source en DOS txns con committed_at distintos (t0 y t1) → 3 filas en t0, 2 en t1.
  • Setear datasets.source_watermark del output a t0 → la ventana (t0, high] cubre SOLO las 2 filas de t1.
  • Correr el compute delta (snapshot_replace incremental) → upsert nativo de las 2 filas de t1.
  • AserciĂłn: el output refleja el merge de SOLO la ventana (las 2 de t1 mergeadas por PK), no las 5 completas. Paridad contra el mismo deltaPredicate corrido en PG.
  • Repetir el gate windowAllowsDelta: un SNAPSHOT en la ventana → sourceDelta=undefined → cae a replace (verificar que NO intenta delta nativo).

DoD B4: delta incremental nativo verde E2E (upsert/append windowed) + el gate windowAllowsDelta respetado + la asunciĂłn dataset_rows-no-drenado documentada + port-migration.md/2d3-tier-b-approach.md marcan B4.


7. Decisiones abiertas (con recomendaciĂłn)

DecisiĂłnRecomendaciĂłn
¿Nativizar el READ del delta ahora?No — opción A (read PG) es correcta mientras nada se drene. B/C al llegar Fase 5 / Karma.
Guard source-nativo-drenadoDiferir — documentar TODO(Fase 5); activar con el drain, no antes (evita falsos bloqueos en la ventana canary).
¿Write mode de snapshot_replace + delta?upsert (el delta es aditivo → merge por PK). El probe#0 decide replace (1er build) vs upsert (re-deploy). Ya cableado.
ÂżB (native-scan) o C (Karma) para el SOTA?Karma subsume delta+cursor; B es un stopgap solo si Fase 5 llega antes que Karma. No construir B especulativamente.

8. Resumen ejecutivo

B4 = delta-aware computeReplaceSql (3 estrategias) + quitar 1 throw. El WRITE ya está (B1/B2/B3); el delta solo narra el compute a la ventana (low,high] del ledger PG. Correcto y seguro hoy (garantía windowAllowsDelta + asunción dataset_rows-legible). El único trabajo real de fondo — nativizar el READ de la ventana — se difiere a Fase 5/Karma con la dependencia documentada, no escondida.

Documento de approach. Al implementar, cruzar con 2d3-tier-b-approach.md §4 (que este doc expande) + paso6-write-approach.md §2e.