Ir al contenido
SparquetSparquet

Un pipeline, muchas ejecuciones

Copiar un pipeline para cambiar un filtro es como un repositorio termina con catorce archivos casi idénticos que se van desviando entre sí. Los parámetros existen para evitar exactamente eso.

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

Un archivo, un camino de código, tres jobs. Corregir un bug lo corrige en todas partes.

Python En el archivo Uso
"br" br un fragmento de ruta, un literal en SQL
42 42 un umbral
True true mantiene un paso protegido por skip_if_false
False (vacío) omite ese paso
["A", "B"] 'A', 'B' un IN (...) de SQL
[1, 2] 1, 2 un IN (...) numérico
[] (vacío) falsy — omite el paso

skip_if_false lee el valor después de la sustitución:

// 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": { } }

Esa tercera forma es lo que reemplaza una pila de archivos casi duplicados: un solo pipeline que cubre varios flujos, con las diferencias expresadas como condiciones en lugar de copias.

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

Los fragmentos se parametrizan como cualquier otra cosa — la sustitución ocurre después de que el include se expande. La ruta es relativa al archivo del pipeline, los includes no se anidan, y la directiva funciona solo en el transformations de nivel superior.

Mantenlos fuera del pipeline y fuera de tu 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",
})

En Databricks, los parámetros de job y los widgets se mapean naturalmente sobre este diccionario. En Airflow, vienen de la configuración de la ejecución del DAG. Eso es lo que hace que un archivo sirva para una carga programada, un backfill y una ejecución de depuración.

Studio lista los parámetros de un Job y renderiza un input para cada uno — tipado como texto, número, switch o chips — y los envía con una ejecución local. También detecta con el linter los dos desajustes que muerden: un {param} usado pero nunca declarado, y un parámetro declarado que nada usa.

  • Nombra por significado, no por posición: {region}, no {p1}.
  • Documéntalos en la description del pipeline — el archivo es la documentación.
  • No parametrices la forma. Los parámetros son valores. Cuando dos variantes necesitan estructuras distintas, skip_if_false es la herramienta; cuando necesitan pipelines distintos, escribe dos.
  • Nunca pongas un secreto en un parámetro que termine en el commit. Lee las credenciales del entorno en el job que lanza la ejecución, y pasa aquí solo valores no sensibles.