API de Python y CLI
Sparquet
Sección titulada «Sparquet»El punto de entrada. Posee la sesión de Spark y comparte los engines de transformación y validación entre todas las ejecuciones.
from sparquet import Sparquet
fw = Sparquet(spark={"app_name": "daily_load", "master": "local[*]"})
fw.run("pipelines/orders.json")fw.run("pipelines/customers.json")
fw.stop()Crear una instancia por proceso y ejecutar muchos pipelines a través de ella es la forma prevista: la sesión se construye una vez, y las vistas temporales que escribe un pipeline son visibles para el siguiente.
run(config_path, input_df=None, columns=None, params=None)
Sección titulada «run(config_path, input_df=None, columns=None, params=None)»Ejecuta un pipeline descrito por un archivo.
| Argumento | Tipo | Propósito |
|---|---|---|
config_path |
str |
ruta al archivo JSON |
input_df |
DataFrame |
reemplaza la lectura del input — el pipeline parte de este DataFrame |
columns |
dict |
columnas literales inyectadas antes de las transformaciones |
params |
dict |
sustituciones de {param} |
result = fw.run( "orders.json", params={"region": "br", "since": "2026-01-01"}, columns={"load_id": "2026-01-15T03:00:00Z"},)columns agrega cada entrada como una columna literal (F.lit(value)) justo después de la lectura — la forma limpia de estampar un id de lote o un timestamp de ejecución sin editar el archivo.
run_from_dict(config, input_df=None, columns=None, params=None)
Sección titulada «run_from_dict(config, input_df=None, columns=None, params=None)»Igual, a partir de un diccionario ya en memoria:
result = fw.run_from_dict({ "name": "inline", "input": {"format": "delta", "path": "sales.orders"}, "output": {"format": "parquet", "path": "/tmp/orders"},})Útil cuando un pipeline se genera — por un notebook, un test o el runner local de Studio.
Registro
Sección titulada «Registro»fw.register_reader("my_format", MyReader)fw.register_writer("my_format", MyWriter)fw.register_transformation("normalize_text", NormalizeText)fw.register_validator("no_future_date", NoFutureDateValidator)Consulta Extensión.
Detiene la sesión de Spark subyacente. En Databricks, donde la sesión es de la plataforma, déjala corriendo.
Pipeline
Sección titulada «Pipeline»Un pipeline, sin el envoltorio del framework — útil en tests y cuando gestionas la sesión tú mismo.
from sparquet import Pipeline
result = Pipeline.from_file("orders.json").run()result = Pipeline.from_dict({...}).run()Pipeline también acepta engines inyectados, que es como pruebas un pipeline con una transformación personalizada registrada solo para ese caso.
PipelineResult
Sección titulada «PipelineResult»Cada ejecución devuelve este objeto. Nunca lanza excepciones — las fallas vuelven como datos.
@dataclassclass PipelineResult: pipeline_name: str success: bool rows_read: int = 0 rows_written: int = 0 validation_results: list[ValidationResult] = [] output_metrics: list[OutputMetrics] = [] error: str | None = None output_df: DataFrame | None = None skipped: bool = False
def summary(self) -> str: ...| Campo | Significado |
|---|---|
success |
False cuando la ejecución falló; revisa error para saber por qué |
skipped |
True cuando stop_if_empty terminó la ejecución — success sigue siendo True |
rows_read |
filas leídas del input (0 cuando se inyectó un DataFrame) |
rows_written |
total de filas escritas — la suma de los conteos por destino |
validation_results |
una entrada por regla: rule_type, passed, message, failed_count, severity, metric_value, check_name |
output_metrics |
una OutputMetrics por destino (ver abajo) |
output_df |
el DataFrame transformado, disponible cuando se inyectó input_df |
error |
el mensaje de la falla |
output_metrics
Sección titulada «output_metrics»Una entrada por destino que el pipeline escribió, en orden:
@dataclassclass OutputMetrics: format: str path: str mode: str rows_written: intCada rows_written se cuenta sobre el DataFrame final de ese destino — después de sus propias transformations y proyección de columnas, justo antes de la escritura — así que es exacto incluso cuando una cadena por destino explota o agrega filas. PipelineResult.rows_written es la suma de estos conteos.
result = fw.run("orders.json")for m in result.output_metrics: print(f"{m.format:8} {m.mode:9} {m.rows_written:>10} {m.path}")result = fw.run("orders.json")
if result.skipped: log.info("nothing to process")elif not result.success: raise RuntimeError(result.error)else: log.info(result.summary()) for check in result.validation_results: if not check.passed: alert(f"{check.rule_type}: {check.message} ({check.failed_count} rows)")Encadenar pipelines en un solo job
Sección titulada «Encadenar pipelines en un solo job»fw = Sparquet(spark={"app_name": "registration"})
for conf in ["conf_a.json", "conf_b.json", "conf_c.json"]: result = fw.run(conf, params={"dt_ref": dt_ref}) if not result.success: raise RuntimeError(f"{conf}: {result.error}")
# every pipeline above wrote the same temp view; this one publishes itfinal = fw.run("conf_commit.json")fw.stop()Cada pipeline escribe una vista temporal de staging, y el último la valida y la publica. Como comparten una sesión, las vistas sobreviven entre llamadas.
python -m sparquet.cli pipeline.jsonO, cuando se instala como paquete, a través del script de consola:
sparquet pipeline.jsonLa CLI parsea el archivo, lo ejecuta e imprime el resultado estructurado. Es el punto de entrada correcto para un scheduler o un contenedor: sin script envoltorio que mantener.
Logging estructurado
Sección titulada «Logging estructurado»Cada línea de log es un objeto JSON en stderr:
{"timestamp": "2026-01-15T03:00:12.918Z", "level": "INFO", "message": "Pipeline finished", "pipeline": "orders_curated", "rows_written": 41233}Ese formato encaja directamente en Datadog, CloudWatch o Splunk sin un parser. El runner local de Studio captura los mismos registros y los muestra en el panel de ejecución.