Tu primer pipeline
Cinco minutos, un archivo, un conjunto de datos real en disco. Sin Studio, sin clúster, sin base de datos.
-
Crea los datos
Sección titulada «Crea los datos»Terminal window mkdir -p sparquet-demo/data && cd sparquet-demodata/pedidos.csv id,cliente,pais,status,cantidad,precio_unitario,fecha_pedido1,Ada,BR,CONFIRMADO,3,25.50,2026-01-042,Linus,PT,CANCELADO,1,80.00,2026-01-053,Grace,BR,CONFIRMADO,2,15.00,2026-01-054,Alan,ES,CONFIRMADO,7,9.90,2026-01-064,Alan,ES,CONFIRMADO,7,9.90,2026-01-06La última fila está duplicada a propósito — el pipeline la eliminará.
-
Escribe el pipeline
Sección titulada «Escribe el pipeline»pedidos.json {"name": "pedidos_curados","description": "Limpia el feed de pedidos y lo escribe como Parquet.","input": {"format": "csv","path": "data/pedidos.csv"},"transformations": [{ "type": "filter", "condition": "status = 'CONFIRMADO'" },{ "type": "cast", "columns": { "cantidad": "int", "precio_unitario": "double", "fecha_pedido": "date" } },{ "type": "with_column", "column": "ingresos", "expression": "cantidad * precio_unitario" },{ "type": "drop_duplicates", "columns": ["id"] }],"validations": {"on_failure": "warn","rules": [{ "type": "not_null", "columns": ["id", "cliente"] },{ "type": "unique", "columns": ["id"] },{ "type": "row_count", "min": 1 }]},"output": {"format": "parquet","path": "out/pedidos","mode": "overwrite","partition_by": ["pais"]}} -
Ejecútalo
Sección titulada «Ejecútalo»run.py from sparquet import Sparquetfw = Sparquet(spark={"app_name": "demo", "master": "local[*]"})result = fw.run("pedidos.json")print(result.summary())for check in result.validation_results:print(f" {check.rule_type}: {'ok' if check.passed else check.message}")fw.stop()Terminal window python run.pyO directamente desde la terminal, sin ningún archivo Python:
Terminal window python -m sparquet.cli pedidos.json -
Lee el resultado
Sección titulada «Lee el resultado»Directoriosparquet-demo
Directoriodata
- pedidos.csv
Directorioout
Directoriopedidos
Directoriopais=BR
- part-0000….parquet
Directoriopais=ES
- part-0000….parquet
- pedidos.json
- run.py
Sobreviven tres filas: el pedido cancelado se filtra y el duplicado se elimina.
ingresosse calculó y la salida está particionada por país.
Qué acaba de pasar
Sección titulada «Qué acaba de pasar»El framework ejecutó el documento de arriba abajo:
| Paso | Qué se ejecutó | Dónde se define |
|---|---|---|
| 1 | Leyó data/pedidos.csv con cabecera y tipos inferidos |
input |
| 2 | Mantuvo los confirmados, ajustó tipos, calculó ingresos, quitó duplicados — en ese orden |
transformations |
| 3 | Midió tres reglas de calidad sin alterar los datos | validations |
| 4 | Escribió Parquet particionado por pais |
output |
Dos detalles que conviene interiorizar ya:
- Las transformaciones se ejecutan en el orden en que las escribes. Filtrar después de un
selectque ya descartó la columna falla — la lista es una secuencia, no un conjunto. - Las validaciones informan, no limpian.
not_nullcuenta nulos; nunca elimina filas. Para eso estáfilter.
Cambia una cosa
Sección titulada «Cambia una cosa»Prueba cada una sobre el archivo que acabas de escribir:
Escribe en dos sitios — sustituye output por outputs:
"outputs": [ { "format": "parquet", "path": "out/pedidos", "mode": "overwrite", "partition_by": ["pais"] }, { "format": "csv", "path": "out/informe", "mode": "overwrite", "columns": ["id", "cliente", "ingresos"] }]Convierte el país en un parámetro:
{ "type": "filter", "condition": "pais = '{pais}'" }fw.run("pedidos.json", params={"pais": "BR"})Falla la ejecución con datos malos — cambia on_failure a "fail" y observa cómo el pipeline se detiene antes de escribir nada.
Siguiente
Sección titulada «Siguiente»- Instalación — extras, drivers y el Studio.
- Transformaciones — las veinte, campo a campo.
- Guías — recetas de principio a fin.