Dónde corren los pipelines
Sparquet es una biblioteca. Donde corre PySpark, corre — y el gestor de sesión se adapta al entorno que encuentra.
| Entorno | Sesión | El bloque spark |
|---|---|---|
| Databricks | reutiliza la sesión activa | ignorado |
| EMR / Dataproc / Synapse | creada con tus configs |
configs aplicados, master ignorado |
| Local | creada con tus configs |
totalmente aplicado, incluido master |
La sesión es un singleton por proceso: el primer Sparquet la crea y cada pipeline posterior la comparte. Eso es lo que hace barato correr muchos pipelines en un solo job, y lo que les permite pasarse datos entre sí por temp views.
from sparquet import Sparquet
fw = Sparquet(spark={"app_name": "dev", "master": "local[*]"})print(fw.run("pipelines/orders.json").summary())fw.stop()Bueno para desarrollo y pruebas. En Windows, Spark necesita winutils.exe y HADOOP_HOME antes de tocar el filesystem local.
%pip install sparquetdbutils.library.restartPython()from sparquet import Sparquet
fw = Sparquet() # reutiliza la sesión del notebook/jobresult = fw.run("/Workspace/Repos/team/pipelines/orders.json", params={"since": dbutils.widgets.get("since")}) # NO llames fw.stop() aquíGuarda los pipelines en un Repo para versionarlos con el resto del código. Los parámetros del job mapean directo a params. Los secretos van en un scope, expuestos como variables de entorno al job.
spark-submit \ --py-files pipelines.zip \ run_job.py --config pipelines/orders.json --since 2026-01-01import argparsefrom sparquet import Sparquet
parser = argparse.ArgumentParser()parser.add_argument("--config", required=True)parser.add_argument("--since", required=True)args = parser.parse_args()
fw = Sparquet(spark={"app_name": "orders"})result = fw.run(args.config, params={"since": args.since})fw.stop()
raise SystemExit(0 if result.success else 1)Salir con código distinto de cero ante un fallo es lo que permite que el step del cluster y tu scheduler se enteren.
FROM apache/spark-py:v3.5.0USER rootRUN pip install --no-cache-dir sparquetCOPY pipelines/ /opt/pipelines/USER sparkENTRYPOINT ["python", "-m", "sparquet.cli"]docker run --rm my-image /opt/pipelines/orders.jsonEmpaquetar pipelines
Sección titulada «Empaquetar pipelines»Los archivos JSON son código — versiónalos con todo lo demás.
repo/├── pipelines/│ ├── orders.json│ ├── customers.json│ └── shared/│ └── standard_filters.json├── jobs/│ └── run_daily.py└── tests/ └── test_pipelines.pyComo un pipeline es dato, es testeable sin cluster: parsea cada archivo en CI y falla en los que no, luego corre los pequeños contra datos de ejemplo.
import json, pathlibfrom sparquet.core.config import PipelineConfig
def test_every_pipeline_parses(): for path in pathlib.Path("pipelines").glob("*.json"): PipelineConfig.from_dict(json.loads(path.read_text(encoding="utf-8")))Agendamiento
Sección titulada «Agendamiento»Sparquet corre un pipeline; no decide cuándo. Cualquier scheduler funciona, porque el punto de entrada es un proceso Python normal:
python -m sparquet.cli pipelines/orders.json# AirflowPythonOperator( task_id="orders", python_callable=lambda **ctx: run_pipeline( "pipelines/orders.json", since=ctx["data_interval_start"].isoformat() ),)Ramifica según el resultado en vez de las excepciones — PipelineResult nunca lanza:
result = fw.run(config, params=params)if result.skipped: return "nothing to do"if not result.success: raise AirflowFailException(result.error)Dependencias por pipeline
Sección titulada «Dependencias por pipeline»Los paquetes de conector se declaran en el archivo, lo que mantiene el job auto-descriptivo:
{ "spark": { "configs": { "spark.jars.packages": "io.delta:delta-spark_2.12:3.2.0" } } }En un cluster compartido, instalar el paquete como biblioteca de cluster es más rápido — no se resuelve en cada ejecución.
Logging
Sección titulada «Logging»Cada línea de log es un objeto JSON en stderr, que cae directo en Datadog, CloudWatch o Splunk sin parser:
{"timestamp":"2026-01-15T03:00:12.918Z","level":"INFO","message":"Pipeline finished","pipeline":"orders_curated","rows_written":41233}Alerta en level: ERROR y en los fallos de validación de la tabla de report — esos dos cubren los modos de fallo que importan.