Seu primeiro pipeline
Cinco minutos, um arquivo, um dataset real em disco. Sem Studio, sem cluster, sem banco.
-
Crie os dados
Seção intitulada “Crie os dados”Terminal window mkdir -p sparquet-demo/data && cd sparquet-demodata/pedidos.csv id,cliente,pais,status,quantidade,preco_unitario,data_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-06A última linha está duplicada de propósito — o pipeline vai removê-la.
-
Escreva o pipeline
Seção intitulada “Escreva o pipeline”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"]}} -
Execute
Seção intitulada “Execute”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.pyOu direto do terminal, sem arquivo Python nenhum:
Terminal window python -m sparquet.cli pedidos.json -
Leia o resultado
Seção intitulada “Leia o resultado”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.
receitafoi calculada e a saída está particionada por país.
O que aconteceu
Seção intitulada “O que aconteceu”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
selectque descartou a coluna quebra — a lista é uma sequência, não um conjunto. - Validações reportam, não limpam.
not_nullconta nulos; nunca remove linhas. Para remover, usefilter.
Mude uma coisa
Seção intitulada “Mude uma coisa”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.
Próximo
Seção intitulada “Próximo”- Seu primeiro Job no Studio — o mesmo pipeline, desenhado.
- Conceitos — o modelo de execução por trás do que você acabou de rodar.
- Transformações — todas as vinte, campo a campo.