Pular para o conteúdo
SparquetSparquet

O JSON do pipeline

Um documento descreve um pipeline. Esta página é o schema inteiro; cada seção liga para as páginas de detalhe.

{
"name": "string", // required
"description": "string", // optional
"spark": { // optional
"app_name": "string",
"master": "string",
"configs": { "spark.sql.shuffle.partitions": "200" }
},
"input": { // required
"format": "csv|parquet|delta|iceberg|txt|view",
"path": "string",
"options": {}
},
"transformations": [], // optional, ordered
"validations": { // optional
"on_failure": "fail|warn|skip",
"report": { "format": "csv", "path": "/dq/report", "mode": "append" },
"rules": []
},
"output": { // required — or "outputs"
"format": "parquet",
"path": "string",
"mode": "overwrite|append|merge|ignore|error",
"partition_by": [],
"columns": [],
"options": {},
"transformations": []
}
}

Obrigatório. Identifica o pipeline nos logs, no relatório de validação e em PipelineResult.pipeline_name. Use um slug estável e amigável ao sistema de arquivos: orders_curated.

Texto opcional. Ele viaja com o arquivo e aparece no Studio, o que o torna a documentação mais barata que um pipeline pode carregar.

Configurações opcionais de sessão.

Chave Default Notas
app_name Sparquet Sobrescrito pelo construtor Sparquet(spark=…)
master local[*] Aplicado apenas no ambiente local
configs {} Qualquer configuração spark.*, aplicada quando a sessão é criada

O uso real comum é declarar pacotes de conectores:

{ "spark": { "configs": { "spark.jars.packages": "io.delta:delta-spark_2.12:3.2.0" } } }

Obrigatório. A fonte única do pipeline.

Chave Tipo Obrigatório
format string sim — convertido para minúsculas no parse
path string sim — tabela, caminho, view ou tópico dependendo do formato
options object não — passado ao reader
{ "input": { "format": "delta", "path": "sales.orders", "options": { "versionAsOf": "12" } } }

Todo formato legível está listado em Conectores. Fontes adicionais entram por meio das transformações join e union.

Após a leitura, o framework adiciona uma coluna ingestion_ts com o timestamp da leitura. Remova-a com drop ou um select se o schema do destino for fixo.

Opcional, ordenado. Veja Transformações para todos os vinte tipos, e Parâmetros para skip_if_false e placeholders.

Chaves desconhecidas dentro de uma transformação são mantidas, não rejeitadas — o que é o que permite às transformações customizadas registradas viajarem pelo mesmo arquivo.

Bloco de qualidade opcional, medido após cada transformação e antes de qualquer escrita.

{
"validations": {
"on_failure": "warn",
"report": { "format": "csv", "path": "/dq/orders", "mode": "append" },
"rules": [
{ "type": "not_null", "columns": ["id"] },
{ "type": "unique", "columns": ["id"] }
]
}
}

Veja Validações para cada regra e o que cada modo on_failure faz.

Um deles é obrigatório. Um único objeto, ou uma lista de destinos.

Chave Tipo Default
format string
path string
mode string overwrite
partition_by lista de colunas []
columns lista de colunas todas as colunas
options object {}
transformations list []
{
"outputs": [
{ "format": "delta", "path": "analytics.orders", "mode": "overwrite" },
{ "format": "csv", "path": "/exports/summary", "columns": ["id", "revenue"] }
]
}

Quando tanto output quanto outputs estão presentes, outputs vence e output é ignorado. Veja Saídas para transformações por destino, projeções e semântica de merge.

  • format é convertido para minúsculas tanto na entrada quanto na saída; mode não é, e alguns writers o comparam de forma sensível a maiúsculas.
  • Chaves desconhecidas sobrevivem. Qualquer coisa que o schema não conheça é carregada para os params da transformação ou ignorada, nunca rejeitada — compatibilidade futura para registros customizados.
  • A substituição de {param} acontece antes do parse, no texto bruto. Um placeholder pode portanto aparecer em qualquer lugar, inclusive dentro de um número ou de uma chave.
  • $include é expandido antes do parse, e apenas no transformations de nível superior.
{
"name": "customer_revenue",
"description": "Daily revenue per customer, curated and published to the lakehouse.",
"input": {
"format": "delta",
"path": "sales.orders"
},
"transformations": [
{ "type": "filter", "condition": "status = 'CONFIRMED' AND ordered_at >= '{since}'" },
{ "type": "stop_if_empty", "message": "No confirmed orders since {since}" },
{ "type": "with_column", "column": "revenue", "expression": "quantity * unit_price" },
{ "type": "checkpoint" }
],
"validations": {
"on_failure": "warn",
"report": { "format": "delta", "path": "quality.pipeline_report", "mode": "append" },
"rules": [
{ "type": "not_null", "columns": ["id", "customer_id"] },
{ "type": "unique", "columns": ["id"] },
{ "type": "row_count", "min": 1 }
]
},
"outputs": [
{
"format": "delta",
"path": "analytics.orders",
"mode": "overwrite",
"partition_by": ["ordered_at"]
},
{
"format": "delta",
"path": "analytics.customer_revenue",
"mode": "merge",
"transformations": [
{
"type": "group_by",
"by": ["customer_id"],
"agg": ["sum(revenue) as revenue", "count(*) as orders"]
}
],
"options": { "merge_keys": ["customer_id"] }
}
]
}
fw.run("customer_revenue.json", params={"since": "2026-01-01"})