Ir al contenido
SparquetSparquet

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.

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.py

Como 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, pathlib
from 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")))

Sparquet corre un pipeline; no decide cuándo. Cualquier scheduler funciona, porque el punto de entrada es un proceso Python normal:

Terminal window
python -m sparquet.cli pipelines/orders.json
# Airflow
PythonOperator(
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)

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.

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.