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 aCOMPLETED, 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 reescribir — commitLakehouseFiles → /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é hace | Registra 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 |
| Coste | Cero reescritura | Una reescritura (ml-runner-side, eficiente) |
Identidad (__row_index/__row_id) | ❌ no | ✅ sí |
| Dataset resultante | Degradado (sin keyset; has_identity mentiría) | Primera clase, idéntico a cualquier dataset Carbon |
| Fichero | Referencia externa (Parquet del partner) | Parquet propio bajo la tabla |
| Código nuevo | mí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). Reutilizarun_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/synclee 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 eldataDestinationa 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.ts—landConsumedTransfer(row): inspecciona schema →createDatasetArtifact→initLakehouseTable→withNativeControlPlane(callback = ingestParquet).- Hook en el worker (
sync-worker.ts, pasoCOMPLETED):COMPLETED → LANDING → LANDED(best-effort), guardatarget_dataset_id. consume: fijardataDestinational 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):
| Pieza | Ruta | Nota |
|---|---|---|
Endpoint /lakehouse/ingest-parquet | services/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 Node | lib/dataspaces/lakehouse-landing.ts | landConsumedTransfer; usa workerDb (PoolClient) para createDatasetArtifact. |
| Cliente del endpoint | lib/dataspaces/edc-client.ts o lib/lakehouse/ | thin fetch a /lakehouse/ingest-parquet (patrón ingest-client.ts). |
| Hook worker | editar lib/workers/edc/sync-worker.ts | en COMPLETED → llamar landConsumedTransfer. |
| Destino en consume | editar app/api/dataspaces/consume/route.ts | dataDestination = inbox lakehouse + keyName único. |
Reutilizado TAL CUAL (cero cambios al core):
| Reúso | Ruta |
|---|---|
| Control plane (ledger + sync_log + stats + schemaVersion) | withNativeControlPlane() · lib/lakehouse/iceberg-native-write.ts:105 |
| Crear tabla Iceberg | initLakehouseTable() · lib/lakehouse/fanout-client.ts:80 |
| Crear dataset + artifact | createDatasetArtifact(client, input) · lib/workers/pipelines/dataset-output-materialiser.ts:587 |
| BASE → Iceberg schema (SSOT) | resolveIcebergSchema() · lib/types/iceberg-types.ts |
| Chunked writer / add_files internos | services/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
- Empírico: confirmar el comportamiento del sink con keyName plano único (¿respeta
<transferId>.parquet?). Ajustar la estrategia de descubrimiento. - ml-runner:
/lakehouse/ingest-parquet(leer R2 Parquet → snapshot nativo). Probar aislado con un Parquet en R2. - Node:
lib/dataspaces/lakehouse-landing.ts+ cliente del endpoint. - Wire: hook en el worker (
COMPLETED → landConsumedTransfer → LANDED) +consumefija eldataDestination. - 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-parquetregistra un Parquet de R2 como snapshot Iceberg nativo (con identidad). -
landConsumedTransfercrea el dataset (datasets+project_files) y lo llena víawithNativeControlPlane. - El worker, en
COMPLETED, aterriza y seteatarget_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
| Riesgo | Mitigación |
|---|---|
| Clave/ubicación del objeto aterrizado (bug prefijo Fase 2) | keyName plano único por transfer + verificación por listing (§3). |
| Schema mismatch Parquet↔tabla | Inferir del Parquet (pyarrow→BASE) dentro de ingest-parquet; validar nombres. |
has_identity incorrecto | Estrategia B reescribe con identidad → has_identity=true es correcto (por eso B, no A). |
| Landing falla tras transfer OK | Best-effort + estado error; reintentable sin re-transferir (el Parquet sigue en R2). |
| Proyecto/carpeta destino | Resolver un proyecto "Dataspaces" por workspace; detalle menor. |
| Coste de reescritura en datasets grandes | ml-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
| Concepto | Ruta | Líneas |
|---|---|---|
withNativeControlPlane() | lib/lakehouse/iceberg-native-write.ts | 105-226 |
initLakehouseTable / commitLakehouseFiles | lib/lakehouse/fanout-client.ts | 80-186 |
ingestDatasetToIceberg (patrón cliente) | lib/lakehouse/ingest-client.ts | 50-79 |
commit_add_files (add_files) | services/ml-runner/app/lakehouse/writer.py | 475-516 |
| Endpoints lakehouse | services/ml-runner/app/routers/lakehouse.py | — |
createDatasetArtifact | lib/workers/pipelines/dataset-output-materialiser.ts | 587-609 |
| Patrón de referencia (sync-worker) | lib/workers/iceberg/sync-worker.ts | 113-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.