Ir al contenido
SparquetSparquet

Parámetros y variables

Cuatro mecanismos convierten un documento estático en uno reutilizable. Se resuelven en momentos distintos, y saber cuál es cuál evita la mayor parte de la confusión.

Mecanismo Se resuelve Proviene de
{param} antes de parsear el argumento params
skip_if_false al aplicar una transformación un valor tras la sustitución
{{variable}} durante la ejecución una transformación collect
$include antes de parsear otro archivo

Los marcadores {param} se reemplazan en el texto crudo del archivo, antes de que se convierta en JSON. Pueden aparecer en cualquier parte — dentro de un string, un número o incluso una clave.

{
"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"],
})

El tipo de Python decide cómo se escribe el valor en el archivo:

Python Se convierte en Uso típico
str / int / float str(value) una ruta, un nombre, un número
True "true" mantiene un paso (con skip_if_false)
False "" (vacío) omite un paso
["a", "b"] 'a', 'b' una cláusula SQL IN (...)
[1, 2] 1, 2 un IN (...) numérico
[] "" (vacío) falsy — omite el paso

Un marcador sin clave correspondiente queda literal en el archivo. Eso es intencional: permite que un fragmento lleve parámetros opcionales sin fallar cuando están ausentes.

Cualquier transformación acepta esta metaclave. Tras la sustitución, el engine decide:

Valor tras la sustitución Resultado
"" (vacío) omitido — desde False, una lista vacía o un param ausente
una expresión que evalúa a un booleano omitido cuando es false
cualquier otro valor no vacío se ejecuta
// 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": { } }

La expresión ve solo literales — los valores ya sustituidos — nunca columnas del DataFrame. Es un interruptor por ejecución, no una condición a nivel de fila. Una expresión que no logra parsearse se trata como “no omitir”, así que un error de tipeo ejecuta el paso en lugar de descartarlo silenciosamente.

{{variable}} se resuelve durante la ejecución, con valores que el propio pipeline calculó. El patrón para el que existe: empujar una lista de claves a una lectura posterior para que la fuente pueda saltar datos.

[
{ "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 es la forma declarativa de df.select(col).distinct().collect() seguido de isin(...) — el truco que permite que Delta y Parquet salten archivos en lugar de escanear la tabla.

  • collect ejecuta distinct().collect() en el driver, así que haz checkpoint primero: sobre un plan no materializado se recalcula todo el linaje.
  • El formateo coincide con el SQL que necesitas: los strings se vuelven 'a', 'b' (comillas escapadas), los números 1, 2.
  • Una colección vacía se renderiza como NULL, así que IN (NULL) no casa con nada — el comportamiento correcto cuando el conjunto de trabajo está vacío.
  • El store se comparte con las with_transformations anidadas, así que las variables recolectadas en la cadena externa son visibles dentro del lado derecho de un join.
  • Se limpia al inicio de cada run(), así que nada se filtra entre pipelines.
  • Una {{variable}} que aún no existe queda literal en lugar de lanzar una excepción.

$include inserta un fragmento en las transformations de nivel 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"] }
]
  • La ruta es relativa al archivo del pipeline.
  • El fragmento es un único objeto o una lista.
  • La sustitución de {param} ocurre después de que el include se expande, así que los fragmentos compartidos pueden parametrizarse.
  • Los includes anidados no se expanden, y la directiva solo funciona en el arreglo transformations de nivel 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
})

Un archivo, un camino de código, cuatro jobs diferentes según lo que pases.