Ir al contenido
SparquetSparquet

Salidas

Un pipeline escribe un destino con output, o varios con outputs. Cada destino se configura de forma independiente, y todos parten del mismo DataFrame transformado.

{
"outputs": [
{ "format": "delta", "path": "analytics.orders", "mode": "overwrite" },
{ "format": "csv", "path": "/exports/summary", "columns": ["id", "revenue"] }
]
}
Clave Tipo Default Propósito
format string conector, en minúsculas al parsear
path string ruta, tabla, vista o tópico
mode string overwrite comportamiento de escritura
partition_by list [] particionamiento físico, solo formatos de archivo
columns list todas proyección de columnas para este destino
options object {} opciones del conector
transformations list [] reformateo solo para este destino
Modo Comportamiento
overwrite reemplaza el destino
append agrega filas
merge upsert — Delta e Iceberg

merge requiere options.merge_keys:

{
"format": "delta",
"path": "analytics.orders",
"mode": "merge",
"options": {
"merge_keys": ["order_id"],
"merge_condition": "T.deleted = false"
}
}

T es la tabla destino, S el DataFrame entrante.

columns selecciona lo que este destino escribe, sin tocar los demás:

{
"outputs": [
{ "format": "parquet", "path": "/dw/orders_full" },
{ "format": "csv", "path": "/exports/finance", "columns": ["order_id", "revenue", "country"] }
]
}

Cuando los destinos necesitan formas diferentes — no solo columnas diferentes — dale a cada uno su propia cadena. Se ejecuta después de las validaciones, solo sobre ese destino.

{
"outputs": [
{
"format": "kafka",
"path": "orders-events",
"transformations": [
{ "type": "struct", "column": "payload", "fields": { "id": "order_id", "total": "revenue" } },
{ "type": "with_column", "column": "value", "expression": "to_json(payload)" }
],
"options": { "bootstrap_servers": "broker:9092", "value_column": "value" }
},
{
"format": "delta",
"path": "analytics.customer_revenue",
"mode": "merge",
"transformations": [
{
"type": "group_by",
"by": ["customer_id"],
"agg": ["sum(revenue) as revenue", "count(*) as orders"]
}
],
"options": { "merge_keys": ["customer_id"] }
}
]
}

Todos los tipos de transformación están disponibles aquí, incluyendo join y las {{variables}} de runtime.

Para cada destino, en este orden:

  1. sus propias transformations
  2. la proyección columns
  3. la escritura

Lo que significa que una proyección ve las columnas que produjeron las propias transformaciones del destino, no aquellas con las que terminó la cadena principal.

Las vistas temporales convierten un conjunto de pipelines en un job componible: varios pipelines escriben la misma vista de staging, y uno final la valida y la publica.

// pipeline 1..N
{ "output": { "format": "view", "path": "orders_staging", "mode": "overwrite" } }
// final pipeline
{ "input": { "format": "view", "path": "orders_staging" } }

Todos deben ejecutarse en la misma sesión de Spark — el mismo proceso — que es exactamente lo que hace Sparquet cuando llamas a run() varias veces sobre una misma instancia.

result = fw.run("pipeline.json")
result.rows_written # total across every destination
for m in result.output_metrics:
print(m.format, m.mode, m.path, m.rows_written)

Cada destino reporta sus propias OutputMetrics (format, path, mode, rows_written) en result.output_metrics. El conteo se toma sobre el DataFrame final de ese destino — después de sus propias transformations y proyección de columnas, justo antes de la escritura — así que es exacto incluso cuando una cadena por destino explota o agrega filas.