Published

🤖 2d.3 · B5 — Model-node native write-arm (predicciones → 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 · B5 — Model-node native write-arm (predicciones → target Iceberg)

Entregable de approach (2026-07-09). El último caso de 2d3-tier-b-approach.md: el write nativo de un nodo model (predicciones). Sus DOS sink.run (model-node-materialiser.ts:249 per-window + :327 truncate-on-reuse) lanzan deferred to 2d.3. Grounded en model-node-materialiser.ts + predictionOutput.ts.

Por qué merece espectro: a diferencia de B1-B4 (donde el compute es una proyección SQL), el model-node tiene dos asimetrías propias que no aparecen en el deploy-materialiser — la CORRELACIÓN y la ACUMULACIÓN. Separarlas es el diseño.

✅ IMPLEMENTADO (correlación A + acumulación A, 2026-07-09) — canary §5e verde E2E. writeRowsToTrident (nuevo en native-output-write.ts) + el loop bufferea las ventanas en la rama iceberg del sink.run → UN solo replace tras el loop + __source_row_index__ como columna explícita del schema nativo + truncate-on-reuse native → no-op. Verificado: 2 ventanas (srcKeys [10,20,30],[40,50]) → 5 filas en un snapshot (acumulación) + __source_row_index__ preservado (correlación). Divergencia de schema (native lleva la columna, PG el row_index) documentada como opt-in; unificación = follow-up cuando native sea default.


1. Cómo escribe el model-node hoy (grounded)

materialiseModelNode (Track A, windowed, cap MAX_ROWS=50 000):

  1. ensureProducedDataset — reuse → DELETE FROM dataset_rows (truncate) [sink.run #2]; else create.
  2. Loop por ventanas de WINDOW_SIZE=1000 (keyset por row_index >= nextStart):
    • lee la ventana proyectando SOURCE_ROW_KEY (= row_index del source) + los inputs mapeados;
    • predice la ventana (ml-runner);
    • escribe: row_index = srcKey (el row_index del source), data = predicción [sink.run #1];
    • avanza nextStart = last srcKey + 1.
  3. Finaliza datasets (schema = solo columnas predichas, row_count).

Dos hechos que definen B5:

  • Correlación: SOURCE_ROW_KEY = '__source_row_index__' (predictionOutput.ts:32). El output NO lleva columna de key — la correlación predicción↔fila-fuente cabalga el row_index (el output reusa el row_index del source). El schema mostrado (buildModelOutputSchema) = solo las columnas predichas.
  • Acumulación: el write es por-ventana (N sink.run, todas con el mismo transactionId). En PG cada ventana INSERTA (acumula en dataset_rows).

2. Las dos asimetrías con el nativo

2a. CORRELACIÓN — row_index lo posee ml-runner

En un write nativo ml-runner estampa __row_index secuencial (0,1,2… en orden de escritura) → el output NO puede fijar row_index = srcKey. La correlación se pierde. Peor: ya es frágil bajo el linchpin — el hydrate-to-temp de una fuente nativa reasigna row_index sintético (source-resolver.ts:91), así que un model encadenado sobre un upstream nativo ya perdería la correlación por row_index.

⇒ La correlación robusta = __source_row_index__ como columna explícita (no como row_index).

2b. ACUMULACIÓN — replace por-ventana borraría todo menos la última

Un facade.write({replace}) por ventana sobrescribe el snapshot anterior → solo sobreviviría la última ventana. Un append por ventana acumula pero prolifera snapshots (N por deploy). ⇒ hace falta acumular las ventanas y hacer UN solo replace al final (el model capea a 50k filas → el buffer en memoria es viable).


3. Espectro de opciones

Correlación:

OpciónQuéTrade-off
A · columna nativa (recomendada)native carga __source_row_index__ como columna explícita (en data + datasets.schema); PG intacto (row_index=srcKey)schema nativo diverge del PG (una columna extra) — pero es HONESTO (native no puede reusar row_index) y opt-in (canary)
B · unificarañadir __source_row_index__ a AMBOS (buildModelOutputSchema)consistente, pero cambia el shape del output PG + downstream
C · diferir (orden de fila)no carga key; ml-runner estampa 0,1,2 en orden srcKey asccorrelación solo si el source es contiguo 0..N-1; frágil con huecos

Acumulación:

OpciónQuéTrade-off
A · buffer + single replace (recomendada)acumular todas las ventanas en memoria → 1 facade.write({replace})O(rows) memoria, acotado por MAX_ROWS=50k (ya es el cap del model)
B · append por-ventana1ª replace, resto appendprolifera snapshots (N por deploy)
C · stagingbuffer→1 snapshot server-sidecomplejidad de ingest.stream; innecesario a 50k

4. El cut recomendado (correlación A + acumulación A)

  1. Resolver el sink UNA vez (ya se hace). Ramas por sink.sink:
    • postgres → el loop per-window INSERTA como hoy (intacto).
    • iceberg → el loop bufferea cada ventana con {...predicción, [SOURCE_ROW_KEY]: srcKey}; el write se hace UNA vez tras el loop.
  2. sink.run #2 (truncate-on-reuse): native → no-op (el replace final subsume el truncate). PG → DELETE como hoy.
  3. Tras el loop, si native: writeRowsToTrident({ datasetId, columns: nativeSchema, rows: buffer, writeMode: 'replace', sourceRefs:[{source}] }) — nuevo helper (hermano en-memoria de writeOutputToTrident, que es SQL-compute). nativeSchema = outputSchema + { name: '__source_row_index__', type: 'Integer', nullable: false }.
  4. Finalize: datasets.schema = nativeSchema cuando native (la columna de correlación es declarada → round-trip + los models encadenados la leen); row_count = written (= filas del buffer = CommitResult.rows).

Atomicidad: withNativeControlPlane es all-or-nothing por snapshot; el model corre dentro del SAVEPOINT del caller. Un fallo a mitad → nada committeado (ni PG-metadata ni snapshot). El replace único (no per-window) hace el write nativo naturalmente atómico.


5. La divergencia de schema (el wart honesto)

Un predictions dataset NATIVO lleva __source_row_index__; el PG no (usa row_index). Es análogo a __row_id/__row_index (identidad que se materializa distinto por almacenamiento). Documentado, opt-in (defaultMode=pg). El canvas (buildModelOutputSchema) muestra solo las columnas predichas — la de correlación es infraestructura, no una columna de negocio. Unificación (opción B) = follow-up cuando native sea el default: mover la correlación a __source_row_index__ en AMBOS y retirar el reuse de row_index.


6. Cambios de código (ordenados)

  1. native-output-write.ts: writeRowsToTrident (rows en memoria → facade.write; cap NATIVE_WRITE_MAX_ROWS). Hermano de writeOutputToTrident.
  2. model-node-materialiser.ts:
    • loop: rama sink.sink==='iceberg' → bufferea {...pred, [SOURCE_ROW_KEY]: srcKey}; PG intacto.
    • sink.run #2: native → no-op.
    • tras el loop: native → writeRowsToTrident({replace}) con nativeSchema.
    • finalize: schema/row_count native-aware.
  3. Canary §5e (verificación).

7. Verificación — canary §5e

El canary no puede correr una predicción real (necesita un model deployado + ml-runner /predict). Prueba el mecanismo (que es lo que B5 añade):

  • Output de model fresco (native, main.test), schema [{prediction, Double}, {__source_row_index__, Integer}].
  • Simular 2 ventanas (srcKeys [10,20,30] y [40,50]) → un buffer acumulado → 1 writeRowsToTrident({replace}).
  • Read-back: 5 filas (acumulación: todas las ventanas, no solo la última) + cada __source_row_index__ preservado (10..50) → la correlación sobrevive el round-trip nativo.

DoD B5: §5e verde (acumulación + correlación por columna) + los 2 sink.run ya no lanzan (native cablea) + datasets.schema native lleva __source_row_index__ + docs marcan B5 → Tier B CERRADO.


8. Resumen ejecutivo

B5 = writeRowsToTrident (nuevo) + bufferear el loop + __source_row_index__ como columna nativa + truncate→no-op. Las dos asimetrías propias del model-node (ml-runner posee row_index → correlación por columna; write per-window → buffer + single replace) son todo el diseño; el resto reusa el substrato nativo. La divergencia de schema (native lleva la columna de correlación, PG el row_index) es honesta y opt-in; unificar es follow-up cuando native sea default.

Documento de approach. Cierra el espectro de Tier B (B1-B5). Cruzar con 2d3-tier-b-approach.md §5.