Pular para o conteúdo
SparquetSparquet

Saídas

Um pipeline escreve um destino com output, ou vários com outputs. Cada destino é configurado de forma independente, e todos partem do mesmo DataFrame transformado.

{
"outputs": [
{ "format": "delta", "path": "analytics.orders", "mode": "overwrite" },
{ "format": "csv", "path": "/exports/summary", "columns": ["id", "revenue"] }
]
}
Chave Tipo Default Propósito
format string conector, convertido para minúsculas no parse
path string caminho, tabela, view ou tópico
mode string overwrite comportamento de escrita
partition_by list [] particionamento físico, apenas formatos de arquivo
columns list todas projeção de colunas para este destino
options object {} opções do conector
transformations list [] remodelagem apenas para este destino
Modo Comportamento
overwrite substitui o destino
append adiciona linhas
merge upsert — Delta e Iceberg

merge requer options.merge_keys:

{
"format": "delta",
"path": "analytics.orders",
"mode": "merge",
"options": {
"merge_keys": ["order_id"],
"merge_condition": "T.deleted = false"
}
}

T é a tabela destino, S o DataFrame de entrada.

columns seleciona o que este destino escreve, sem afetar os demais:

{
"outputs": [
{ "format": "parquet", "path": "/dw/orders_full" },
{ "format": "csv", "path": "/exports/finance", "columns": ["order_id", "revenue", "country"] }
]
}

Quando os destinos precisam de formas diferentes — não só colunas diferentes — dê a cada um sua própria cadeia. Ela roda após as validações, apenas naquele destino.

{
"outputs": [
{
"format": "kafka",
"path": "orders-events",
"transformations": [
{ "type": "struct", "column": "payload", "fields": { "id": "order_id", "total": "revenue" } },
{ "type": "with_column", "column": "value", "expression": "to_json(payload)" }
],
"options": { "bootstrap_servers": "broker:9092", "value_column": "value" }
},
{
"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"] }
}
]
}

Todos os tipos de transformação estão disponíveis aqui, incluindo join e {{variables}} de runtime.

Para cada destino, nesta ordem:

  1. suas próprias transformations
  2. a projeção de columns
  3. a escrita

O que significa que uma projeção enxerga as colunas que as próprias transformações do destino produziram, não aquelas com que a cadeia principal terminou.

Temp views transformam um conjunto de pipelines em um job componível: vários pipelines escrevem na mesma view de staging, e um final valida e a publica.

// pipeline 1..N
{ "output": { "format": "view", "path": "orders_staging", "mode": "overwrite" } }
// final pipeline
{ "input": { "format": "view", "path": "orders_staging" } }

Todos eles precisam rodar na mesma sessão Spark — o mesmo processo — que é exatamente o que o Sparquet faz quando você chama run() várias vezes em uma instância.

result = fw.run("pipeline.json")
result.rows_written # total across every destination
for m in result.output_metrics:
print(m.format, m.mode, m.path, m.rows_written)

Cada destino reporta seu próprio OutputMetrics (format, path, mode, rows_written) em result.output_metrics. A contagem é feita no DataFrame final daquele destino — após suas próprias transformations e projeção de colunas, logo antes da escrita — então é exata mesmo quando uma cadeia por destino explode ou agrega linhas.