✍️ Paso 6 — Encaminar la escritura al Data Space (el tridente)
Approach fundado en el código (2026-07-03). Objetivo: que toda la estructura de escritura de Carbon deje de aterrizar en Postgres
dataset_rowsy escriba directamente al tridente — R2 (bytes) + Lakekeeper (punteros) + PG (metadata/gobernanza) — vía la facadewrite()(loadItemForConsumption).Companion del linchpin (READ) ya cerrado: dataspace-router.md · source-of-truth.md · tracker vivo: port-migration.md.
Sintetizado de 3 auditorías sobre el terreno (motor de write nativo · re-arquitectura del pipeline write · paisaje completo de writers). Karma NO es prerequisito — ml-runner ya escribe Iceberg hoy; Karma es el reemplazo posterior que elimina el PG del medio.
1. El motor de write nativo — production-ready
La puerta ya existe y está verificada E2E:
facade.write({ kind:'rows', rows: AsyncIterable, columns, writeMode, sourceRefs, metadata })
→ writeIcebergNativeSnapshot (lib/lakehouse/iceberg-native-write.ts)
→ ingestDatasetToIceberg (NDJSON → ml-runner POST /lakehouse/ingest)
· ml-runner ESTAMPA la identidad: __row_index (0-based), __row_id (uuid), __created_at
· escribe Parquet a R2 + registra la tabla Iceberg (Lakekeeper)
→ withNativeControlPlane: ledger(open) → ingest → commit → iceberg_sync_log(has_identity)
→ datasets stats → dataset_schema_versions (fail-loud, sin fallback PG)
Hechos clave:
- La identidad la posee ml-runner — el productor NO reasigna
row_index; solo streamea las filas de negocio. (Disuelve el "problema más duro" que la auditoría del pipeline write señaló:row_number()en PG deja de ser necesario.) - El control plane es atómico y fail-loud (nunca cae a PG → nunca split-brain).
- Patrón de referencia:
DatasetWriter.writeDatasetIcebergNative(flatten batches → per-row stream →writeIcebergNativeSnapshot). Los writers Tier-1 ya lo usan. facade.write()es correcta — el throw en PG-sink es el guard fail-loud, no un bug. Paso 6 no cambia la facade; secuencia por-writer.
2. El paisaje de writers — un stack de 3 tiers
| Tier | Writers | maxMode | Estado |
|---|---|---|---|
| 1 · nativo (hecho) | sync.full, sync.incremental, sync.append, cdc.append, ingest.api | iceberg_native | R2 por defecto. Referencia. |
| 2 · flippable | manual.table · pipeline.output | iceberg_native / pg | manual.table = flip (hecho) · pipeline.output = proyecto acoplado |
| 3 · diferido (falta primitiva Iceberg) | ingest.stream (no streaming open/commit → proliferación de snapshots) · row.edit (no point-DELETE en PyIceberg) | pg | → Karma / staging-append. Clamp PG. |
| 4 · PG permanente | dataset.manifest | pg | control-plane (punteros de fichero, no dato tabular) |
maxMode es el techo de capacidad (writer-flags.ts): el router clampa cualquier flag a él → un Tier-3 no puede flipear por accidente aunque un flag lo pida.
3. El reframe — la escritura no puede adelantar a la lectura
Igual que el linchpin READ tuvo que rutarse antes de que el pipeline pudiera flipear (split-brain de input, ya resuelto), un output no puede irse a nativo hasta que sus LECTORES estén facade-routed — si no, leen un dataset_rows drenado (vacío). Couplings downstream rankeados:
| # | Coupling | Severidad | Qué es |
|---|---|---|---|
| AUDITADO (2026-07-03, 3 agentes) — los 3 couplings "de miedo" ya estaban resueltos por trabajo previo: |
| # | Coupling | Estado real (auditado) |
|---|---|---|
| 1 | objects_resolved / objects grid | ✅ RESUELTO — la vista ya NO hace LEFT JOIN dataset_rows; lee objects.base_properties denormalizado (Track A de objects.sync, mig. 20261228). Seguro para datasets nativos, sin migración. |
| 2 | ledger de delta (pipeline_build_transactions) | ✅ RESUELTO — storage-agnostic (solo metadata txn/committed_at). El downstream delta corre en cualquier substrato. |
| 3 | downstream pipelines leyendo el output | ✅ RESUELTO — el linchpin (source-resolver) hidrata un source nativo. (native→native es un perf cliff, optimización de Fase 4, no bloqueo.) |
| 4 | readers analytics — "gate inerte" (llaman al seam pero hardcodean FROM dataset_rows, ignoran el source) | ✅ CERRADO (2c, 2026-07-04) — graph.kuzu + dashboards.aggregate REFERENCE (2026-07-03) · datasets.export.legacy + model.schema.infer (facade read/stream, 2026-07-04) · sql-editor + dashboards.aggregate QUERY (live-SQL → hydrate-to-inline-jsonb, 2026-07-04). Ningún reader analytics lee ya un dataset_rows drenado. |
Reorientación: el gate más temido (objects grid) ya era seguro. El 2c real = readers analytics, todos migrados. Notas de implementación: kuzu lee filas a memoria →
handle.read()directo; dashboards-reference native →stream()+ agregación en memoria;datasets.export.legacy→stream();model.schema.infer→read()+ descriptorrowCount. Los CTE-readers de SQL ARBITRARIO (dashboards QUERY, sql-editor) NO caben en unread()estructurado → hydrate-to-inline-jsonb (lib/lakehouse/sql-source-hydration.ts): un dataset nativo se hidrata ajsonb_to_recordset(función pura → corre enBEGIN READ ONLY, sinCREATE TEMP TABLEni relajar el sandbox — mejor que la "hydrate-to-temp read-write" que se temía en Fase 3). Hermano interactivo del linchpinsource-resolver.ts(que sí usa temp read-write porque el pipeline compute ya es read-write).
✅
pipeline.outputYA puede flipear GLOBAL limpio — todos los readers analytics son source-respecting (facade-routed). Ya no queda el gate inerte; el sweep original sí sobre-estimó cobertura conflando "seam-wired" con "source-respecting", pero eso está cerrado.
4. El approach por fases
- 2a — Quick win:
manual.table→ nativo ✅ (2026-07-03). FlipdefaultMode pg→iceberg_native. La ruta ya escribía nativo en sink iceberg (writeIcebergNativeSnapshotreplace post-commit); su consumidor primario (pipelines vía el source-resolver del linchpin) maneja el source nativo. Reversible + pinnable por dataset. - 2b — Fase D (gobernanza), paralelo, barato ✅ (2026-07-03).
mergeUpstreamProvenancecableado en los chokepoints de derivación (createOutputDatasetdel pipeline output +createTransformOutputDataset+createJoinOutputDataset) → el output hereda la procedencia del source (counterparty/edc/upstream lineage) en vez de solosource_dataset_id. La gobernanza deja de evaporarse en la derivación. - 2c — El GATE real (readers analytics) ✅ CERRADO (2026-07-04): los couplings de miedo (objects_resolved, ledger, downstream) ya estaban resueltos. Migrados:
graph.kuzu(kuzuFacadeRead) +dashboards.aggregateREFERENCE (native →stream()+aggregateInMemory) (2026-07-03);datasets.export.legacy+model.schema.infer(facaderead/stream+ descriptorrowCount);sql-editor+dashboards.aggregateQUERY (live-SQL → hydrate-to-inline-jsonb,lib/lakehouse/sql-source-hydration.ts) (2026-07-04). Prerequisito de flipearpipeline.output= SATISFECHO — todos los readers analytics son source-respecting; ya no hace falta canary por-dataset evitándolos. - 2d —
pipeline.outputnativo (el proyecto grande, tras 2c): re-arquitectura escalonada compute (PG) →facade.write({replace})→ ml-runner → Iceberg. Replace-first (ml-runner posee row_index; transform/join/split/pure-replace mapean limpio; diferir append + merge-preservante). → APPROACH DE ESPECTRO COMPLETO (5 auditorías): 2d-pipeline-output-flip.md.- 2d.1 ✅ (P0 + Apply-path): los 5 P0 bloqueantes cerrados (expectations source-respecting, model-node source→linchpin, build-lock, count-ratchet, transform-paths seam) + el brazo iceberg del Apply-path completo (transform/join/split) escribiendo el tridente vía
writeOutputToTrident;maxMode pg→iceberg_native; deploy source-respecting (sourceTablepor el linchpin). - 2d.2 ✅ (deploy replace-safe + canary): el deploy materialiser tiene brazo iceberg REPLACE-safe (
computeReplaceSqlpor estrategia, solo primer-deploy con target vacío) →writeOutputToTrident. Canary VERDE E2E (scripts/dataspaces/pipeline-output-canary.ts,npm run canary:pipeline-output) contra infra viva enmain.test: default-safe + flip + write nativo + control-plane + paridad PG↔Iceberg exacta. defaultMode SIGUE pg (flip opt-in por-dataset). - 2d.3 / Tier B ✅ CERRADO — 2d3-tier-b-approach.md: probe#0 (oráculo, no
count(*) dataset_rows) + B1upsert+ B2append+ B3append_new(anti-join leyendo los PKs del target por la facade) + B4 delta nativo (opción A: read ventana=PG, write=dispatch; 2d3-b4-delta-native.md) + B5 model-node (buffer+single-replace +__source_row_index__como columna; 2d3-b5-model-node.md). Canary VERDE E2E (§5a-e). Diferido a Karma/Fase 5: read nativo de la ventana delta, unificación correlación model.
- 2d.1 ✅ (P0 + Apply-path): los 5 P0 bloqueantes cerrados (expectations source-respecting, model-node source→linchpin, build-lock, count-ratchet, transform-paths seam) + el brazo iceberg del Apply-path completo (transform/join/split) escribiendo el tridente vía
- 2e — Diferido a Karma:
ingest.stream(staging-append: buffer batches → un snapshot en commit) ·row.edit(point-DELETE: sin primitiva Iceberg → merge-on-read delete / mutations-as-log). Clamp PG, documentado. No forzar Iceberg.
5. Secuenciación (dependencias)
HECHO ├─ Tier 1 (sync/cdc/ingest.api) nativo
├─ 2a manual.table → nativo (flip)
├─ 2b Fase D (propagación de gobernanza en los chokepoints)
├─ 2c readers analytics → facade (kuzu, dashboards ref+QUERY, export.legacy,
│ schema.infer, sql-editor) — el GATE cerrado; ya no hay gate inerte
├─ 2d.1 P0 (×5) + Apply-path nativo (transform/join/split) + maxMode iceberg_native
├─ 2d.2 deploy replace-safe + CANARY VERDE E2E (main.test, paridad PG↔Iceberg)
└─ 2d.3 / Tier B COMPLETO: probe#0 + upsert + append + append_new (anti-join por
facade) + delta nativo + model-node (buffer+single-replace + __source_row_index__)
DIFERIDO (Karma)
├─ 2e ingest.stream (staging-append)
└─ 2e row.edit (point-DELETE)
PERMANENTE
└─ dataset.manifest (PG por diseño)
Fase D NO es prerequisito de Paso 6 — es ortogonal (qué metadata llevan los writes, no dónde aterrizan). Corre en paralelo.
6. Decisiones abiertas (con recomendación)
| Decisión | Recomendación |
|---|---|
| ¿Karma para el write? | No es prerequisito. ml-runner ya escribe Iceberg (E2E). Karma subsume el compute+write nativo después (elimina el PG del medio). |
| Atomicidad del staged write (2d) | Write-intent ledger: registrar la intención antes del compute; el write nativo es idempotente (replace) → retry seguro. El compute PG rollback + el write nativo commit dejan de ser una tx única. |
| row_index en el pipeline nativo | Que lo posea ml-runner (como sync.*). Replace = limpio. Append continúa desde max (seguro si los builds por dataset están serializados — ya lo están). |
| Merges que leen el target (append-new, snapshot-replace) | La lectura del target existente debe ir por la facade (el target nativo tiene dataset_rows drenado). Replace-puro + transform/join/split primero; merge-preservante después. |
| ¿objects_resolved antes o después? | Antes (2c). Es el gate crítico; y es reader work que continúa la ola de readers. |
7. Anti-patrones que este approach evita
- ❌ Flipear un output a nativo antes que sus lectores → grid/dashboards vacíos.
- ❌ Mixed-mode delta (output nativo / downstream PG) → pérdida silenciosa.
- ❌ Forzar Iceberg donde falta la primitiva (point-DELETE, streaming commit) → esperar a Karma, clamp PG.
- ❌ Reasignar row_index en el productor → dejar que ml-runner posea la identidad.
- ❌ Que la gobernanza se evapore en la derivación → propagarla en los chokepoints (Fase D).
Documento de approach. La escritura sigue al linchpin READ ya cerrado; el gate es objects_resolved (2c). Mantener al día conforme cada writer aterrice en el tridente.