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.
El documento del pipeline
Sección titulada «El documento del pipeline»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": [ … ]}Orden de ejecución
Sección titulada «Orden de ejecución»Esta secuencia es el modelo mental de todo el framework:
- Sustitución de
{param}— los marcadores se reemplazan en el texto crudo, antes de que el JSON se parsee. - Expansión de
$include— los fragmentos referenciados se insertan inline entransformations. - Parse — el documento se convierte en configuración tipada; las claves desconocidas se conservan, no se rechazan.
- Lectura del input — una fuente, a través del registro de readers. Se agrega automáticamente una columna
ingestion_ts. - Transformaciones — aplicadas en el orden del array, con las variables
{{runtime}}resueltas a medida que cada una se ejecuta. - Validaciones — medidas sobre el DataFrame transformado; aquí se escribe el reporte opcional.
- 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.
Dos tipos de marcador
Sección titulada «Dos tipos de marcador»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.
Cómo un lienzo se convierte en un archivo
Sección titulada «Cómo un lienzo se convierte en un archivo»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 shared2. 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_transformationssource → filter ─┴────────────────────────────────────→ join → sinkTodo lo demás es un campo de un nodo. Las notas nunca compilan; los nodos deshabilitados se omiten.
Dónde se ejecuta un pipeline
Sección titulada «Dónde se ejecuta un pipeline»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.
Modelo de fallos
Sección titulada «Modelo de fallos»PipelineResult nunca lanza excepciones. Una ejecución fallida vuelve como datos:
result = fw.run("pipeline.json")
result.success # False when something went wrongresult.error # the message, when it didresult.skipped # True when stop_if_empty ended the run earlyresult.rows_read # rows read from the inputresult.rows_written # rows in the main DataFrame at write timeresult.validation_results # one entry per ruleEsa forma es lo que permite a un orquestador ramificar según el resultado sin envolver todo en un try.
Siguiente
Sección titulada «Siguiente»- El JSON del pipeline — cada campo, con sus valores por defecto.
- Transformaciones — las veinte.
- Guías — recetas de principio a fin.