🤖 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 DOSsink.run(model-node-materialiser.ts:249per-window +:327truncate-on-reuse) lanzandeferred to 2d.3. Grounded enmodel-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 ennative-output-write.ts) + el loop bufferea las ventanas en la ramaicebergdelsink.run→ UN soloreplacetras 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 elrow_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):
ensureProducedDataset— reuse →DELETE FROM dataset_rows(truncate) [sink.run #2]; else create.- Loop por ventanas de
WINDOW_SIZE=1000(keyset porrow_index >= nextStart):- lee la ventana proyectando
SOURCE_ROW_KEY(=row_indexdel source) + los inputs mapeados; - predice la ventana (ml-runner);
- escribe:
row_index = srcKey(elrow_indexdel source),data = predicción[sink.run #1]; - avanza
nextStart = last srcKey + 1.
- lee la ventana proyectando
- 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 elrow_index(el output reusa elrow_indexdel source). El schema mostrado (buildModelOutputSchema) = solo las columnas predichas. - Acumulación: el write es por-ventana (N
sink.run, todas con el mismotransactionId). En PG cada ventana INSERTA (acumula endataset_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ón | Qué | 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 · unificar | añ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 asc | correlación solo si el source es contiguo 0..N-1; frágil con huecos |
Acumulación:
| Opción | Qué | 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-ventana | 1ª replace, resto append | prolifera snapshots (N por deploy) |
| C · staging | buffer→1 snapshot server-side | complejidad de ingest.stream; innecesario a 50k |
4. El cut recomendado (correlación A + acumulación A)
- 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.
sink.run #2(truncate-on-reuse): native → no-op (elreplacefinal subsume el truncate). PG →DELETEcomo hoy.- Tras el loop, si native:
writeRowsToTrident({ datasetId, columns: nativeSchema, rows: buffer, writeMode: 'replace', sourceRefs:[{source}] })— nuevo helper (hermano en-memoria dewriteOutputToTrident, que es SQL-compute).nativeSchema = outputSchema + { name: '__source_row_index__', type: 'Integer', nullable: false }. - Finalize:
datasets.schema = nativeSchemacuando 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)
native-output-write.ts:writeRowsToTrident(rows en memoria →facade.write; capNATIVE_WRITE_MAX_ROWS). Hermano dewriteOutputToTrident.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})connativeSchema. - finalize:
schema/row_countnative-aware.
- loop: rama
- 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 → 1writeRowsToTrident({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.