Ir al contenido
SparquetSparquet

Tu primer pipeline

Cinco minutos, un archivo, un conjunto de datos real en disco. Sin Studio, sin clúster, sin base de datos.

  1. Terminal window
    mkdir -p sparquet-demo/data && cd sparquet-demo
    data/pedidos.csv
    id,cliente,pais,status,cantidad,precio_unitario,fecha_pedido
    1,Ada,BR,CONFIRMADO,3,25.50,2026-01-04
    2,Linus,PT,CANCELADO,1,80.00,2026-01-05
    3,Grace,BR,CONFIRMADO,2,15.00,2026-01-05
    4,Alan,ES,CONFIRMADO,7,9.90,2026-01-06
    4,Alan,ES,CONFIRMADO,7,9.90,2026-01-06

    La última fila está duplicada a propósito — el pipeline la eliminará.

  2. 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"]
    }
    }
  3. run.py
    from sparquet import Sparquet
    fw = 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.py

    O directamente desde la terminal, sin ningún archivo Python:

    Terminal window
    python -m sparquet.cli pedidos.json
    • 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. ingresos se calculó y la salida está particionada por país.

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 select que ya descartó la columna falla — la lista es una secuencia, no un conjunto.
  • Las validaciones informan, no limpian. not_null cuenta nulos; nunca elimina filas. Para eso está filter.

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.