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.
%pip install sparquetdbutils.library.restartPython()from sparquet import Sparquet
fw = Sparquet() # reusa a sessão do notebook/jobresult = fw.run("/Workspace/Repos/team/pipelines/orders.json", params={"since": dbutils.widgets.get("since")}) # NÃO chame fw.stop() aquiGuarde os pipelines num Repo para versioná-los junto com o resto do código. Parâmetros do job mapeiam direto em params. Segredos vão num scope, expostos como variáveis de ambiente para o 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)Sair com código não-zero em caso de falha é o que deixa o step do cluster e o seu scheduler perceberem.
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.jsonEmpacotando pipelines
Seção intitulada “Empacotando pipelines”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.pyComo 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, 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")))Agendamento
Seção intitulada “Agendamento”O Sparquet roda um pipeline; ele não decide quando. Qualquer scheduler funciona, porque o ponto de entrada é um processo 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() ),)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)Dependências por pipeline
Seção intitulada “Dependências por pipeline”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.
Logging
Seção intitulada “Logging”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.