Ir al contenido
SparquetSparquet

Conceptos fundamentales

Todo en Sparquet se deriva de un único orden de ejecución. Apréndelo una vez y el resto de la referencia se lee sola.

Un pipeline tiene cinco partes. Solo tres son obligatorias.

{
"name": "orders_curated", // required
"description": "", // optional
"spark": { "configs": {} }, // optional
"input": { }, // required
"transformations": [ ], // optional, ordered
"validations": { }, // optional
"output": { } // required — or "outputs": [ … ]
}

Esta secuencia es el modelo mental de todo el framework:

  1. Sustitución de {param} — los marcadores se reemplazan en el texto crudo, antes de que el JSON se parsee.
  2. Expansión de $include — los fragmentos referenciados se insertan inline en transformations.
  3. Parse — el documento se convierte en configuración tipada; las claves desconocidas se conservan, no se rechazan.
  4. Lectura del input — una fuente, a través del registro de readers. Se agrega automáticamente una columna ingestion_ts.
  5. Transformaciones — aplicadas en el orden del array, con las variables {{runtime}} resueltas a medida que cada una se ejecuta.
  6. Validaciones — medidas sobre el DataFrame transformado; aquí se escribe el reporte opcional.
  7. Para cada destino — sus propias transformations, luego su proyección de columnas, luego la escritura.

Dos consecuencias con las que la gente tropieza:

  • Las transformaciones se ejecutan antes que las validaciones. Una regla siempre ve los datos ya limpios, nunca la fuente cruda.
  • Las transformaciones por destino se ejecutan después de las validaciones. Reformar los datos para una salida nunca afecta lo que las reglas midieron, ni lo que otra salida recibe.

Las transformaciones cambian los datos. Las validaciones reportan sobre ellos.

Sección titulada «Las transformaciones cambian los datos. Las validaciones reportan sobre ellos.»
Quieres… Usa
Eliminar filas con un id nulo filter en transformations
Saber cuántos ids nulos llegaron not_null en validations
Descartar duplicados drop_duplicates
Hacer fallar la ejecución cuando existen duplicados unique con on_failure: "fail"

Una validación nunca modifica el DataFrame. Esa separación es lo que hace confiable a un reporte de calidad: describe los datos que realmente escribiste.

Sparquet tiene dos mecanismos de sustitución, y se resuelven en momentos distintos.

{param} {{variable}}
Se resuelve antes del parse durante la ejecución
Proviene de el argumento params de la ejecución una transformación collect
Uso típico entorno, fecha, feature flags una lista de claves empujada hacia una lectura posterior
Si no se resuelve queda literal, sin error queda literal, se resuelve después si aparece
// {param}: known when you launch the run
{ "type": "filter", "condition": "region = '{region}'" }
// {{variable}}: computed by the pipeline itself
{ "type": "collect", "column": "customer_id", "as": "active" },
{ "type": "join",
"with": { "format": "delta", "path": "sales.events" },
"with_transformations": [
{ "type": "filter", "condition": "customer_id IN ({{active}})" }
],
"on": "customer_id" }

Consulta Parámetros y variables para las reglas de formato de cada tipo.

Studio no almacena un formato propietario — compila el grafo. Dos reglas definen la correspondencia:

1. La cadena compartida es el pipeline principal. Las transformaciones que todos los destinos tienen en común se convierten en las transformations de nivel superior. Una vez que el grafo se bifurca, cada rama se convierte en las transformations propias de ese destino.

source → filter → cast ─┬─→ [group_by] → delta ← group_by belongs to this output
└─→ parquet ← filter and cast are shared

2. Un join o union toma su segunda fuente de su segundo input. La cadena que alimenta ese handle se convierte en with_transformations.

┌── delta(events) → filter → select ──┐ ← with + with_transformations
source → filter ─┴────────────────────────────────────→ join → sink

Todo lo demás es un campo de un nodo. Las notas nunca compilan; los nodos deshabilitados se omiten.

El framework detecta su entorno y adapta la sesión:

Entorno Comportamiento
Databricks Reutiliza la sesión activa; el bloque spark se ignora
EMR / Dataproc / Synapse Construye una sesión con tus configs
Local Aplica también master (local[*] por defecto)

La sesión es un singleton a nivel de proceso: la primera instancia de Sparquet gana, y los pipelines posteriores la comparten. Eso es lo que hace barato ejecutar varios pipelines en un mismo job.

PipelineResult nunca lanza excepciones. Una ejecución fallida vuelve como datos:

result = fw.run("pipeline.json")
result.success # False when something went wrong
result.error # the message, when it did
result.skipped # True when stop_if_empty ended the run early
result.rows_read # rows read from the input
result.rows_written # rows in the main DataFrame at write time
result.validation_results # one entry per rule

Esa forma es lo que permite a un orquestador ramificar según el resultado sin envolver todo en un try.