Pular para o conteúdo
SparquetSparquet

Onde os pipelines rodam

Sparquet é uma biblioteca. Onde o PySpark roda, ele roda — e o gerenciador de sessão se adapta ao ambiente que encontra.

Ambiente Sessão O bloco spark
Databricks reusa a sessão ativa ignorado
EMR / Dataproc / Synapse criada com seus configs configs aplicados, master ignorado
Local criada com seus configs totalmente aplicado, incluindo master

A sessão é um singleton por processo: o primeiro Sparquet a cria e todo pipeline seguinte a compartilha. É isso que torna barato rodar muitos pipelines num job só, e o que permite que eles passem dados entre si por temp views.

from sparquet import Sparquet
fw = Sparquet(spark={"app_name": "dev", "master": "local[*]"})
print(fw.run("pipelines/orders.json").summary())
fw.stop()

Bom para desenvolvimento e testes. No Windows, o Spark precisa de winutils.exe e HADOOP_HOME antes de tocar no filesystem local.

Os arquivos JSON são código — versione-os com todo o resto.

repo/
├── pipelines/
│ ├── orders.json
│ ├── customers.json
│ └── shared/
│ └── standard_filters.json
├── jobs/
│ └── run_daily.py
└── tests/
└── test_pipelines.py

Como um pipeline é dado, ele é testável sem cluster: faça o parse de todo arquivo no CI e falhe nos que não passarem, depois rode os pequenos contra dados de exemplo.

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")))

O Sparquet roda um pipeline; ele não decide quando. Qualquer scheduler funciona, porque o ponto de entrada é um processo 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()
),
)

Ramifique pelo resultado em vez de por exceções — PipelineResult nunca lança:

result = fw.run(config, params=params)
if result.skipped:
return "nothing to do"
if not result.success:
raise AirflowFailException(result.error)

Pacotes de conector são declarados no arquivo, o que mantém o job auto-descritivo:

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

Num cluster compartilhado, instalar o pacote como biblioteca de cluster é mais rápido — não é resolvido a cada execução.

Toda linha de log é um objeto JSON no stderr, que cai direto em Datadog, CloudWatch ou Splunk sem parser:

{"timestamp":"2026-01-15T03:00:12.918Z","level":"INFO","message":"Pipeline finished","pipeline":"orders_curated","rows_written":41233}

Alerte em level: ERROR e nas falhas de validação da tabela de report — esses dois cobrem os modos de falha que importam.