Published

Fase 4 — Aterrizar el dato consumido en el lakehouse (Iceberg gobernado)

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

Fase 4 — Aterrizar el dato consumido en el lakehouse (Iceberg gobernado)

Documento padre: README.md · Base: fase-3-findings.md (Fase 3 ✅; el worker deja un hook en COMPLETED)
Estado: approach / alcance de la iteración.
Naturaleza: cuando un transfer EDC llega a COMPLETED, el dato consumido es un Parquet en R2. Fase 4 lo convierte en un dataset de Carbon gobernado (tabla Iceberg + ledger + lineaje + visible en el workspace), reutilizando el control plane del lakehouse sin tocar su core.


1. El problema central (y el insight)

El transfer deja un Parquet suelto en R2 (bucket destino). Para que sea un dataset de primera clase de Carbon hay que: crear el artifact (datasets + project_files), crear la tabla Iceberg, y registrar los datos con el control plane (dataset_transactions ledger + iceberg_sync_log + stats).

Insight de la auditoría (clave para el alcance): existe add_files (registrar un Parquet ya en R2 sin reescribircommitLakehouseFiles/lakehouse/commit). Es tentador por zero-copy. PERO el Parquet del partner no trae las columnas de identidad de Carbon (__row_index, __row_id) que las escrituras nativas sí ponen y de las que depende iceberg_sync_log.has_identity=true (keyset pagination, read-router). Un add_files directo de un Parquet ajeno produciría un dataset degradado (sin identidad) y marcaría has_identity incorrectamente.

→ La decisión de alcance gira sobre identidad vs zero-copy.


2. Decisión de encaje — dos estrategias

A. add_files (zero-copy)B. ingest-parquet (nativo)
Qué haceRegistra el Parquet del partner tal cual (commitLakehouseFiles)ml-runner lee el Parquet de R2 y reescribe un snapshot nativo (con identidad), reutilizando el writer chunked
CosteCero reescrituraUna reescritura (ml-runner-side, eficiente)
Identidad (__row_index/__row_id)❌ no✅ sí
Dataset resultanteDegradado (sin keyset; has_identity mentiría)Primera clase, idéntico a cualquier dataset Carbon
FicheroReferencia externa (Parquet del partner)Parquet propio bajo la tabla
Código nuevomínimo (reusa commitLakehouseFiles)1 endpoint ml-runner (aditivo) + orquestación Node

Recomendación: B. El valor de Carbon es tener datasets gobernados y consistentes (lineaje, paginación, read-router). Un dataset "de segunda" (sin identidad) rompería esa promesa.

Naturaleza de B — NO es un stack de ingesta nuevo. Es una fuente nueva para el writer existente. Tu pipeline ya ingesta por fuentes hacia el mismo escritor:
/lakehouse/sync = fuente Postgres, /lakehouse/ingest = fuente NDJSON, los polling-adapters = fuentes DB / S3-object. B añade la fuente "Parquet ya en R2"hermana de /sync (server-side: ml-runner lee y escribe, sin round-trip a Node). Reutiliza run_lakehouse_ingest() tal cual (identidad __row_index/__row_id, chunking, add_files, commit, catálogo, control plane): lo único nuevo es leer el Parquet de R2 y emitir sus filas a esa función.

¿Por qué server-side y no un polling-adapter Node? Porque el s3-adapter difiere Parquet (s3-adapter.ts:20) y Node no lee Parquet — igual que /sync lee Postgres en ml-runner, la fuente Parquet vive en ml-runner. Cero código nuevo aguas abajo; A queda como fast-path zero-copy opcional (ligado al zero-copy real de Fase 5).


3. Sub-problema: ¿dónde/con qué clave aterriza el Parquet?

Hallazgo Fase 2: el sink S3 de EDC ignoró el prefijo del keyName (dataspace-inbox/sample.parquet → aterrizó sample.parquet en la raíz del bucket). Para saber qué objeto produjo cada transfer:

  • En consume: fijar el dataDestination a un bucket de inbox de dataspaces (en el R2 del lakehouse) con un keyName plano y único por transfer (<transferId>.parquet, sin "/") → esquiva el bug del prefijo.
  • En el worker (al COMPLETED): verificar por listing el objeto real (robusto si el sink cambia la clave), y pasar esa ruta R2 a ingest-parquet.

(Comprobar empíricamente el comportamiento exacto del sink con keyName plano — igual que hicimos con v4 — es parte de la iteración.)


4. Alcance de la iteración

Dentro:

  • Endpoint ml-runner /lakehouse/ingest-parquet: dado {dataset_id, namespace, source_path (R2), schema?} → lee el Parquet, infiere/valida schema, reescribe snapshot nativo (identidad + chunked add_files internos), devuelve {identifier, rows, snapshot_id, size_bytes}.
  • lib/dataspaces/lakehouse-landing.tslandConsumedTransfer(row): inspecciona schema → createDatasetArtifactinitLakehouseTablewithNativeControlPlane(callback = ingestParquet).
  • Hook en el worker (sync-worker.ts, paso COMPLETED): COMPLETED → LANDING → LANDED (best-effort), guarda target_dataset_id.
  • consume: fijar dataDestination al inbox del lakehouse con keyName único.

Fuera (a propósito):

  • ❌ Zero-copy real / DataSource Lakekeeper (compartir tabla sin mover) → Fase 5.
  • ❌ Append/incremental/upsert del mismo asset en el tiempo → primer corte = replace one-shot (cada transfer = un snapshot).
  • ❌ UI de la pestaña Dataspaces → Fase 7.
  • ❌ Mapear el dataset consumido a la ontología/objetos → fases de ontología.

5. Piezas y reúso

A crear (poco, aditivo):

PiezaRutaNota
Endpoint /lakehouse/ingest-parquetservices/ml-runner/app/routers/lakehouse.py + lakehouse/parquet_source.py (nuevo, pequeño)Reutiliza run_lakehouse_ingest() tal cual (el mismo writer que /ingest). Nuevo = leer el Parquet de R2 (PyArrow S3FS, path-style) → emitir filas + inferir fields del schema del Parquet.
Orquestación Nodelib/dataspaces/lakehouse-landing.tslandConsumedTransfer; usa workerDb (PoolClient) para createDatasetArtifact.
Cliente del endpointlib/dataspaces/edc-client.ts o lib/lakehouse/thin fetch a /lakehouse/ingest-parquet (patrón ingest-client.ts).
Hook workereditar lib/workers/edc/sync-worker.tsen COMPLETED → llamar landConsumedTransfer.
Destino en consumeeditar app/api/dataspaces/consume/route.tsdataDestination = inbox lakehouse + keyName único.

Reutilizado TAL CUAL (cero cambios al core):

ReúsoRuta
Control plane (ledger + sync_log + stats + schemaVersion)withNativeControlPlane() · lib/lakehouse/iceberg-native-write.ts:105
Crear tabla IceberginitLakehouseTable() · lib/lakehouse/fanout-client.ts:80
Crear dataset + artifactcreateDatasetArtifact(client, input) · lib/workers/pipelines/dataset-output-materialiser.ts:587
BASE → Iceberg schema (SSOT)resolveIcebergSchema() · lib/types/iceberg-types.ts
Chunked writer / add_files internosservices/ml-runner/app/lakehouse/writer.py

6. Contratos concretos

/lakehouse/ingest-parquet (nuevo):

POST  { dataset_id, namespace, source_path: "s3://…/<transferId>.parquet", schema?: BASEField[] }
  →   { identifier, rows, snapshot_id, location, size_bytes }

(Si schema viene, se valida/usa; si no, se infiere del Parquet — pyarrow → BASE.)

landConsumedTransfer(row: DataspaceTransfer) (nuevo, Node):

1. sourcePath ← descubrir el objeto real del transfer (row.config.dest + listing)
2. schema     ← inspeccionar (endpoint o dentro de ingest-parquet)
3. { datasetId } ← createDatasetArtifact(workerDbClient, { workspaceId: row.workspace_id,
                      projectId, name: `EDC: ${row.asset_id}`, service: 'edc_consumer',
                      schema, metadata: { edc_transfer_id: row.edc_transfer_id, counterparty: row.counterparty }})
4. initLakehouseTable(datasetId, resolveIcebergSchema(schema), { writeMode: 'replace' })
5. withNativeControlPlane(
     { datasetId, workspaceId: row.workspace_id, columns: schema, writeMode: 'replace', namespace,
       metadata: { source: 'edc_consumer', edc_transfer_id: row.edc_transfer_id } },
     async () => ingestParquet(datasetId, namespace, sourcePath),  // → NativeIngestOutcome
   )
6. UPDATE dataspace_transfers SET target_dataset_id = datasetId, state = 'LANDED'

Estados (columna state, sin migración nueva): COMPLETED → LANDING → LANDED (o COMPLETED + error si el landing falla — el transfer siguió OK; el landing se reintenta). target_dataset_id ya existe.

Proyecto destino: para el primer corte, un proyecto/carpeta "Dataspaces" por workspace (o el default del workspace). Resolver projectId es un detalle menor (reutilizar el resolver de proyecto por defecto).


7. Runbook

  1. Empírico: confirmar el comportamiento del sink con keyName plano único (¿respeta <transferId>.parquet?). Ajustar la estrategia de descubrimiento.
  2. ml-runner: /lakehouse/ingest-parquet (leer R2 Parquet → snapshot nativo). Probar aislado con un Parquet en R2.
  3. Node: lib/dataspaces/lakehouse-landing.ts + cliente del endpoint.
  4. Wire: hook en el worker (COMPLETED → landConsumedTransfer → LANDED) + consume fija el dataDestination.
  5. E2E: consume → worker → transfer COMPLETED → aparece un dataset en el workspace con las filas del Parquet, gobernado (ledger + sync_log + stats). Verificar con un smoke (patrón de los de Fase 3) contra el conector + Supabase real.

8. Definition of Done

  • /lakehouse/ingest-parquet registra un Parquet de R2 como snapshot Iceberg nativo (con identidad).
  • landConsumedTransfer crea el dataset (datasets+project_files) y lo llena vía withNativeControlPlane.
  • El worker, en COMPLETED, aterriza y setea target_dataset_id + state=LANDED (best-effort; error no pierde el transfer).
  • E2E: un consume acaba con un dataset visible y legible en el workspace (read-router lo lee como cualquier dataset).
  • Cero cambios al control plane / ledger / oracle; solo reúso.

9. Riesgos

RiesgoMitigación
Clave/ubicación del objeto aterrizado (bug prefijo Fase 2)keyName plano único por transfer + verificación por listing (§3).
Schema mismatch Parquet↔tablaInferir del Parquet (pyarrow→BASE) dentro de ingest-parquet; validar nombres.
has_identity incorrectoEstrategia B reescribe con identidad → has_identity=true es correcto (por eso B, no A).
Landing falla tras transfer OKBest-effort + estado error; reintentable sin re-transferir (el Parquet sigue en R2).
Proyecto/carpeta destinoResolver un proyecto "Dataspaces" por workspace; detalle menor.
Coste de reescritura en datasets grandesml-runner chunked (bounded RAM); aceptable para el primer corte (replace one-shot).

10. Qué desbloquea

  • El dato consumido por el dataspace es ya un dataset gobernado de Carbon → pipelines, ML, ontología, lineaje inter-org (el transfer queda ligado al dataset por edc_transfer_id/target_dataset_id).
  • Fase 5 (publicar / zero-copy) reutiliza el mismo esqueleto de landing/registro en sentido inverso.
  • Fase 7 (UI): la pestaña Dataspaces ya puede mostrar "consumido → dataset X".

11. Referencias de la auditoría

ConceptoRutaLíneas
withNativeControlPlane()lib/lakehouse/iceberg-native-write.ts105-226
initLakehouseTable / commitLakehouseFileslib/lakehouse/fanout-client.ts80-186
ingestDatasetToIceberg (patrón cliente)lib/lakehouse/ingest-client.ts50-79
commit_add_files (add_files)services/ml-runner/app/lakehouse/writer.py475-516
Endpoints lakehouseservices/ml-runner/app/routers/lakehouse.py
createDatasetArtifactlib/workers/pipelines/dataset-output-materialiser.ts587-609
Patrón de referencia (sync-worker)lib/workers/iceberg/sync-worker.ts113-146

Siguiente acción al ejecutar: Paso 1 del runbook — confirmar empíricamente la clave de aterrizaje del sink, luego el endpoint /lakehouse/ingest-parquet.