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 |
Modos de escrita
Seção intitulada “Modos de escrita”| 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.
Projeção de colunas
Seção intitulada “Projeção de colunas”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"] } ]}Transformações por destino
Seção intitulada “Transformações por destino”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.
Ordem das operações por destino
Seção intitulada “Ordem das operações por destino”Para cada destino, nesta ordem:
- suas próprias
transformations - a projeção de
columns - 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.
Staging entre pipelines
Seção intitulada “Staging entre pipelines”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.
O que retorna
Seção intitulada “O que retorna”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.