Ir al contenido
SparquetSparquet

API de Python y CLI

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.

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.

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.

Cada ejecución devuelve este objeto. Nunca lanza excepciones — las fallas vuelven como datos.

@dataclass
class 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

Una entrada por destino que el pipeline escribió, en orden:

@dataclass
class OutputMetrics:
format: str
path: str
mode: str
rows_written: int

Cada 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)")
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 it
final = 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.

Terminal window
python -m sparquet.cli pipeline.json

O, cuando se instala como paquete, a través del script de consola:

Terminal window
sparquet pipeline.json

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

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.