Published

✍️ Paso 6 — Encaminar la escritura al Data Space (el tridente)

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...

✍️ 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_rows y escriba directamente al tridente — R2 (bytes) + Lakekeeper (punteros) + PG (metadata/gobernanza) — vía la facade write() (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

TierWritersmaxModeEstado
1 · nativo (hecho)sync.full, sync.incremental, sync.append, cdc.append, ingest.apiiceberg_nativeR2 por defecto. Referencia.
2 · flippablemanual.table · pipeline.outputiceberg_native / pgmanual.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 permanentedataset.manifestpgcontrol-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:

#CouplingSeveridadQué es
AUDITADO (2026-07-03, 3 agentes) — los 3 couplings "de miedo" ya estaban resueltos por trabajo previo:
#CouplingEstado real (auditado)
1objects_resolved / objects gridRESUELTO — 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.
2ledger de delta (pipeline_build_transactions)RESUELTO — storage-agnostic (solo metadata txn/committed_at). El downstream delta corre en cualquier substrato.
3downstream pipelines leyendo el outputRESUELTO — el linchpin (source-resolver) hidrata un source nativo. (native→native es un perf cliff, optimización de Fase 4, no bloqueo.)
4readers 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.legacystream(); model.schema.inferread() + descriptor rowCount. Los CTE-readers de SQL ARBITRARIO (dashboards QUERY, sql-editor) NO caben en un read() estructurado → hydrate-to-inline-jsonb (lib/lakehouse/sql-source-hydration.ts): un dataset nativo se hidrata a jsonb_to_recordset (función pura → corre en BEGIN READ ONLY, sin CREATE TEMP TABLE ni relajar el sandbox — mejor que la "hydrate-to-temp read-write" que se temía en Fase 3). Hermano interactivo del linchpin source-resolver.ts (que sí usa temp read-write porque el pipeline compute ya es read-write).

pipeline.output YA 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). Flip defaultMode pg→iceberg_native. La ruta ya escribía nativo en sink iceberg (writeIcebergNativeSnapshot replace 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). mergeUpstreamProvenance cableado en los chokepoints de derivación (createOutputDataset del pipeline output + createTransformOutputDataset + createJoinOutputDataset) → el output hereda la procedencia del source (counterparty/edc/upstream lineage) en vez de solo source_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.aggregate REFERENCE (native → stream() + aggregateInMemory) (2026-07-03); datasets.export.legacy + model.schema.infer (facade read/stream + descriptor rowCount); sql-editor + dashboards.aggregate QUERY (live-SQL → hydrate-to-inline-jsonb, lib/lakehouse/sql-source-hydration.ts) (2026-07-04). Prerequisito de flipear pipeline.output = SATISFECHO — todos los readers analytics son source-respecting; ya no hace falta canary por-dataset evitándolos.
  • 2d — pipeline.output nativo (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 (sourceTable por el linchpin).
    • 2d.2 ✅ (deploy replace-safe + canary): el deploy materialiser tiene brazo iceberg REPLACE-safe (computeReplaceSql por 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 en main.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) + B1 upsert + B2 append + B3 append_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.
  • 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ónRecomendació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 nativoQue 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.