Pular para o conteúdo
SparquetSparquet

Seu primeiro pipeline

Cinco minutos, um arquivo, um dataset real em disco. Sem Studio, sem cluster, sem banco.

  1. Terminal window
    mkdir -p sparquet-demo/data && cd sparquet-demo
    data/pedidos.csv
    id,cliente,pais,status,quantidade,preco_unitario,data_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

    A última linha está duplicada de propósito — o pipeline vai removê-la.

  2. pedidos.json
    {
    "name": "pedidos_curados",
    "description": "Limpa o feed de pedidos e grava em Parquet.",
    "input": {
    "format": "csv",
    "path": "data/pedidos.csv"
    },
    "transformations": [
    { "type": "filter", "condition": "status = 'CONFIRMADO'" },
    { "type": "cast", "columns": { "quantidade": "int", "preco_unitario": "double", "data_pedido": "date" } },
    { "type": "with_column", "column": "receita", "expression": "quantidade * preco_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

    Ou direto do terminal, sem arquivo Python nenhum:

    Terminal window
    python -m sparquet.cli pedidos.json
    • Directorysparquet-demo
      • Directorydata
        • pedidos.csv
      • Directoryout
        • Directorypedidos
          • Directorypais=BR
            • part-0000….parquet
          • Directorypais=ES
            • part-0000….parquet
      • pedidos.json
      • run.py

    Três linhas sobrevivem: o pedido cancelado é filtrado e a duplicata é removida. receita foi calculada e a saída está particionada por país.

O framework executou o documento de cima para baixo:

Passo O que rodou Onde está definido
1 Leu data/pedidos.csv com header e tipos inferidos input
2 Manteve confirmados, ajustou tipos, calculou receita, removeu duplicatas — nessa ordem transformations
3 Mediu três regras de qualidade sem alterar os dados validations
4 Gravou Parquet particionado por pais output

Dois detalhes que valem internalizar agora:

  • Transformações rodam na ordem em que você escreve. Filtrar depois de um select que descartou a coluna quebra — a lista é uma sequência, não um conjunto.
  • Validações reportam, não limpam. not_null conta nulos; nunca remove linhas. Para remover, use filter.

Experimente cada uma no arquivo que você acabou de escrever:

Grave em dois lugares — troque output por outputs:

"outputs": [
{ "format": "parquet", "path": "out/pedidos", "mode": "overwrite", "partition_by": ["pais"] },
{ "format": "csv", "path": "out/relatorio", "mode": "overwrite", "columns": ["id", "cliente", "receita"] }
]

Transforme o país em parâmetro:

{ "type": "filter", "condition": "pais = '{pais}'" }
fw.run("pedidos.json", params={"pais": "BR"})

Falhe a execução com dado ruim — mude on_failure para "fail" e veja o pipeline parar antes de gravar qualquer coisa.