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 |
Modos de escritura
Sección titulada «Modos de escritura»| 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.
Proyección de columnas
Sección titulada «Proyección de columnas»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"] } ]}Transformaciones por destino
Sección titulada «Transformaciones por destino»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.
Orden de operaciones por destino
Sección titulada «Orden de operaciones por destino»Para cada destino, en este orden:
- sus propias
transformations - la proyección
columns - 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.
Staging entre pipelines
Sección titulada «Staging entre pipelines»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.
Lo que devuelve
Sección titulada «Lo que devuelve»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.