Pular para o conteúdo
SparquetSparquet

Parâmetros e variáveis

Quatro mecanismos transformam um documento estático em um reutilizável. Eles são resolvidos em momentos diferentes, e saber qual é qual evita a maior parte da confusão.

Mecanismo Resolve Vem de
{param} antes do parse o argumento params
skip_if_false ao aplicar uma transformação um valor após a substituição
{{variable}} durante a execução uma transformação collect
$include antes do parse outro arquivo

Os placeholders {param} são substituídos no texto bruto do arquivo, antes de virar JSON. Podem aparecer em qualquer lugar — dentro de uma string, um número ou até uma chave.

{
"input": { "format": "delta", "path": "sales.orders_{region}" },
"transformations": [
{ "type": "filter", "condition": "ordered_at >= '{since}' AND status IN ({statuses})" }
]
}
fw.run("orders.json", params={
"region": "br",
"since": "2026-01-01",
"statuses": ["CONFIRMED", "SHIPPED"],
})

O tipo Python decide como o valor é escrito no arquivo:

Python Vira Uso típico
str / int / float str(value) um caminho, um nome, um número
True "true" mantém um passo (com skip_if_false)
False "" (vazio) pula um passo
["a", "b"] 'a', 'b' uma cláusula SQL IN (...)
[1, 2] 1, 2 um IN (...) numérico
[] "" (vazio) falsy — pula o passo

Um placeholder sem chave correspondente fica literal no arquivo. Isso é intencional: permite que um fragmento carregue parâmetros opcionais sem falhar quando eles estão ausentes.

Qualquer transformação aceita essa meta chave. Após a substituição o engine decide:

Valor após substituição Resultado
"" (vazio) pulado — de False, uma lista vazia ou um param ausente
uma expressão que avalia para booleano pulado quando é false
qualquer outro valor não vazio executa
// switch a whole join on and off per run
{ "type": "join", "skip_if_false": "{enrich}", "with": { }, "on": "customer_id" }
// only apply the filter when a region was given
{ "type": "filter", "skip_if_false": "{region}", "condition": "region = '{region}'" }
// branch on a value
{ "type": "struct", "skip_if_false": "'{flow}' in ('ISSUE', 'ISSUE_AND_REGISTER')", "column": "payload", "fields": { } }

A expressão enxerga apenas literais — os valores já substituídos — nunca colunas do DataFrame. É um interruptor por execução, não uma condição a nível de linha. Uma expressão que falha ao ser parseada é tratada como “não pular”, então um erro de digitação executa o passo em vez de descartá-lo silenciosamente.

{{variable}} é resolvida durante a execução, com valores que o próprio pipeline computou. O padrão para o qual ela existe: empurrar uma lista de chaves para uma leitura posterior, para que a fonte possa pular dados.

[
{ "type": "filter", "condition": "status = 'PENDING'" },
{ "type": "checkpoint" },
{ "type": "collect", "column": "customer_id", "as": "pending_customers" },
{
"type": "join",
"with": { "format": "delta", "path": "sales.bronze_events" },
"with_transformations": [
{ "type": "filter", "condition": "customer_id IN ({{pending_customers}})" },
{ "type": "select", "columns": ["customer_id", "segment"] }
],
"on": "customer_id",
"how": "left"
}
]

Esta é a forma declarativa de df.select(col).distinct().collect() seguida de isin(...) — o truque que permite ao Delta e ao Parquet pular arquivos em vez de escanear a tabela.

  • collect roda distinct().collect() no driver, então faça checkpoint primeiro: num plano não materializado toda a linhagem é recomputada.
  • A formatação corresponde ao SQL que você precisa: strings viram 'a', 'b' (aspas escapadas), números 1, 2.
  • Uma coleção vazia é renderizada como NULL, então IN (NULL) não casa nada — o comportamento correto quando o conjunto de trabalho está vazio.
  • O store é compartilhado com as with_transformations aninhadas, então variáveis coletadas na cadeia externa são visíveis dentro do lado direito de um join.
  • Ele é zerado no início de cada run(), então nada vaza entre pipelines.
  • Uma {{variable}} que ainda não existe fica literal em vez de lançar erro.

$include insere inline um fragmento no transformations de nível superior:

{
"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"] }
]
  • O caminho é relativo ao arquivo do pipeline.
  • O fragmento é um único objeto ou uma lista.
  • A substituição de {param} acontece depois que o include é expandido, então fragmentos compartilhados podem ser parametrizados.
  • Includes aninhados não são expandidos, e a diretiva só funciona no array transformations de nível superior.
{
"name": "regional_load",
"input": { "format": "delta", "path": "sales.orders" },
"transformations": [
{ "type": "filter", "condition": "region = '{region}'" },
{ "type": "filter", "skip_if_false": "{products}", "condition": "product_id IN ({products})" },
{ "type": "stop_if_empty", "message": "Nothing for {region}" },
{ "type": "checkpoint" },
{ "type": "collect", "column": "customer_id", "as": "customers" },
{
"type": "join",
"skip_if_false": "{enrich}",
"with": { "format": "delta", "path": "crm.customers" },
"with_transformations": [
{ "type": "filter", "condition": "customer_id IN ({{customers}})" }
],
"on": "customer_id",
"how": "left"
}
],
"output": { "format": "delta", "path": "analytics.orders_{region}", "mode": "overwrite" }
}
fw.run("regional_load.json", params={
"region": "br",
"products": ["P1", "P2"], # [] would skip the second filter entirely
"enrich": True, # False would skip the join
})

Um arquivo, um caminho de código, quatro jobs diferentes dependendo do que você passar.