Ir al contenido
SparquetSparquet

El JSON del pipeline

Un documento describe un pipeline. Esta página es el esquema completo; cada sección enlaza a las páginas de detalle.

{
"name": "string", // required
"description": "string", // optional
"spark": { // optional
"app_name": "string",
"master": "string",
"configs": { "spark.sql.shuffle.partitions": "200" }
},
"input": { // required
"format": "csv|parquet|delta|iceberg|txt|view",
"path": "string",
"options": {}
},
"transformations": [], // optional, ordered
"validations": { // optional
"on_failure": "fail|warn|skip",
"report": { "format": "csv", "path": "/dq/report", "mode": "append" },
"rules": []
},
"output": { // required — or "outputs"
"format": "parquet",
"path": "string",
"mode": "overwrite|append|merge|ignore|error",
"partition_by": [],
"columns": [],
"options": {},
"transformations": []
}
}

Obligatorio. Identifica el pipeline en los logs, en el reporte de validación y en PipelineResult.pipeline_name. Usa un slug estable y amigable con el sistema de archivos: orders_curated.

Prosa opcional. Viaja con el archivo y aparece en Studio, lo que la convierte en la documentación más barata que un pipeline puede llevar.

Configuración de sesión opcional.

Clave Default Notas
app_name Sparquet Sobrescrito por el constructor Sparquet(spark=…)
master local[*] Aplicado solo en el entorno local
configs {} Cualquier ajuste spark.*, aplicado al crear la sesión

El uso real más común es declarar paquetes de conectores:

{ "spark": { "configs": { "spark.jars.packages": "io.delta:delta-spark_2.12:3.2.0" } } }

Obligatorio. La única fuente del pipeline.

Clave Tipo Obligatorio
format string sí — en minúsculas al parsear
path string sí — tabla, ruta, vista o tópico según el formato
options object no — se pasa al lector
{ "input": { "format": "delta", "path": "sales.orders", "options": { "versionAsOf": "12" } } }

Todos los formatos legibles están listados en Conectores. Las fuentes adicionales entran a través de las transformaciones join y union.

Después de la lectura, el framework agrega una columna ingestion_ts con el timestamp de lectura. Descártala con drop o un select si el esquema del destino es fijo.

Opcional, ordenado. Consulta Transformaciones para los veinte tipos, y Parámetros para skip_if_false y los marcadores.

Las claves desconocidas dentro de una transformación se conservan, no se rechazan — que es lo que permite que las transformaciones personalizadas registradas viajen a través del mismo archivo.

Bloque de calidad opcional, medido después de cada transformación y antes de cualquier escritura.

{
"validations": {
"on_failure": "warn",
"report": { "format": "csv", "path": "/dq/orders", "mode": "append" },
"rules": [
{ "type": "not_null", "columns": ["id"] },
{ "type": "unique", "columns": ["id"] }
]
}
}

Consulta Validaciones para cada regla y qué hace cada modo de on_failure.

Uno de ellos es obligatorio. Un único objeto, o una lista de destinos.

Clave Tipo Default
format string
path string
mode string overwrite
partition_by lista de columnas []
columns lista de columnas todas las columnas
options object {}
transformations list []
{
"outputs": [
{ "format": "delta", "path": "analytics.orders", "mode": "overwrite" },
{ "format": "csv", "path": "/exports/summary", "columns": ["id", "revenue"] }
]
}

Cuando output y outputs están ambos presentes, outputs gana y output se ignora. Consulta Salidas para transformaciones por destino, proyecciones y semántica de merge.

  • format se pasa a minúsculas tanto en input como en output; mode no, y algunos escritores lo comparan de forma sensible a mayúsculas.
  • Las claves desconocidas sobreviven. Cualquier cosa que el esquema no conozca se lleva a los params de la transformación o se ignora, nunca se rechaza — compatibilidad hacia adelante para registros personalizados.
  • La sustitución de {param} ocurre antes de parsear, sobre el texto crudo. Por lo tanto, un marcador puede aparecer en cualquier parte, incluso dentro de un número o una clave.
  • $include se expande antes de parsear, y solo en las transformations de nivel superior.
{
"name": "customer_revenue",
"description": "Daily revenue per customer, curated and published to the lakehouse.",
"input": {
"format": "delta",
"path": "sales.orders"
},
"transformations": [
{ "type": "filter", "condition": "status = 'CONFIRMED' AND ordered_at >= '{since}'" },
{ "type": "stop_if_empty", "message": "No confirmed orders since {since}" },
{ "type": "with_column", "column": "revenue", "expression": "quantity * unit_price" },
{ "type": "checkpoint" }
],
"validations": {
"on_failure": "warn",
"report": { "format": "delta", "path": "quality.pipeline_report", "mode": "append" },
"rules": [
{ "type": "not_null", "columns": ["id", "customer_id"] },
{ "type": "unique", "columns": ["id"] },
{ "type": "row_count", "min": 1 }
]
},
"outputs": [
{
"format": "delta",
"path": "analytics.orders",
"mode": "overwrite",
"partition_by": ["ordered_at"]
},
{
"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"] }
}
]
}
fw.run("customer_revenue.json", params={"since": "2026-01-01"})