Pular para o conteúdo
SparquetSparquet

Um pipeline, muitas execuções

Copiar um pipeline para mudar um filtro é como um repositório acaba com quatorze arquivos quase idênticos que divergem entre si. Os parâmetros existem para evitar exatamente isso.

{
"name": "regional_orders",
"input": { "format": "delta", "path": "sales.orders" },
"transformations": [
{ "type": "filter", "condition": "region = '{region}' AND ordered_at >= '{since}'" },
{ "type": "filter", "skip_if_false": "{products}", "condition": "product_id IN ({products})" },
{ "type": "stop_if_empty", "message": "Nothing for {region} since {since}" }
],
"output": { "format": "delta", "path": "analytics.orders_{region}", "mode": "overwrite" }
}
for region in ["br", "pt", "es"]:
fw.run("regional_orders.json", params={
"region": region,
"since": since,
"products": [], # empty → the product filter is skipped entirely
})

Um arquivo, um caminho de código, três jobs. Corrigir um bug o corrige em todos os lugares.

Python No arquivo Uso
"br" br um fragmento de caminho, um literal em SQL
42 42 um threshold
True true mantém um passo protegido por skip_if_false
False (vazio) pula esse passo
["A", "B"] 'A', 'B' um IN (...) de SQL
[1, 2] 1, 2 um IN (...) numérico
[] (vazio) falsy — pula o passo

skip_if_false lê o valor depois da substituição:

// on/off from a boolean parameter
{ "type": "join", "skip_if_false": "{enrich}", "with": { }, "on": "customer_id" }
// only when a value was provided
{ "type": "filter", "skip_if_false": "{region}", "condition": "region = '{region}'" }
// branch on which flow is running
{ "type": "struct", "skip_if_false": "'{flow}' in ('ISSUE', 'ISSUE_AND_REGISTER')", "column": "payload", "fields": { } }

Essa terceira forma é o que substitui uma pilha de arquivos quase duplicados: um pipeline cobrindo vários fluxos, com as diferenças expressas como condições em vez de cópias.

{
"transformations": [
{ "$include": "shared/standard_filters.json" },
{ "type": "with_column", "column": "revenue", "expression": "quantity * unit_price" }
]
}
shared/standard_filters.json
[
{ "type": "filter", "condition": "status = '{status}'" },
{ "type": "drop_duplicates", "columns": ["id"] }
]

Fragmentos são parametrizados como qualquer outra coisa — a substituição acontece depois de o include ser expandido. O caminho é relativo ao arquivo do pipeline, includes não aninham, e a diretiva funciona apenas no transformations de nível superior.

Mantenha-os fora do pipeline e fora do seu código:

import os
from datetime import date
fw.run("regional_orders.json", params={
"region": os.environ["REGION"],
"since": date.today().isoformat(),
"enrich": os.environ.get("ENRICH", "false") == "true",
})

No Databricks, os parâmetros de job e os widgets mapeiam naturalmente para esse dicionário. No Airflow, eles vêm da configuração de execução do DAG. É isso que faz um único arquivo servir a uma carga agendada, a um backfill e a uma execução de depuração.

O Studio lista os parâmetros de um Job e renderiza um campo de entrada para cada um — tipado como texto, número, switch ou chips — e os envia com uma execução local. Ele também aponta os dois descompassos que incomodam: um {param} usado mas nunca declarado, e um parâmetro declarado que nada usa.

  • Nomeie por significado, não por posição: {region}, não {p1}.
  • Documente-os na description do pipeline — o arquivo é a documentação.
  • Não parametrize a forma. Parâmetros são valores. Quando duas variantes precisam de estruturas diferentes, skip_if_false é a ferramenta; quando precisam de pipelines diferentes, escreva dois.
  • Nunca coloque um segredo em um parâmetro que acaba sendo commitado. Leia as credenciais do ambiente no job que dispara a execução, e passe aqui apenas valores não sensíveis.