Ir al contenido
SparquetSparquet

Joins y pushdown en tiempo de ejecución

{
"type": "join",
"with": { "format": "delta", "path": "sales.customers" },
"with_transformations": [
{ "type": "filter", "condition": "active = true" },
{ "type": "select", "columns": ["customer_id", "segment"] },
{ "type": "distinct" }
],
"on": "customer_id",
"how": "left"
}

El lado derecho se lee en línea y se reforma con with_transformations antes del join. Tres hábitos rinden de inmediato:

  • select solo lo que necesitas. Un join arrastra cada columna de ambos lados hacia el shuffle.
  • distinct cuando el lado derecho pueda repetir la clave, o un left join multiplicará filas silenciosamente.
  • Filtra temprano, en el lado derecho, no después del join.
"on": "customer_id" // same name on both sides
"on": ["customer_id", "order_date"] // composite key
"on": "l.customer_id = r.id AND l.dt = r.date" // different names, or an inequality

El DataFrame izquierdo lleva el alias l y el derecho r. Esos alias existen solo en la expresión on y después del join — dentro de with_transformations el lado derecho todavía no tiene alias, así que usa nombres de columna simples ahí.

Valor Conserva
inner (por defecto) filas que coinciden en ambos lados
left cada fila izquierda; columnas derechas nulas cuando no hay coincidencia
right cada fila derecha
full todo de ambos lados
leftsemi filas izquierdas que tienen coincidencia — sin columnas derechas
leftanti filas izquierdas sin coincidencia — la consulta “¿qué falta?”
cross producto cartesiano

leftsemi es el que hay que usar cuando solo necesitas filtrar por existencia: nunca duplica filas y nunca agrega columnas.

El problema: tu conjunto de trabajo es de 40 000 pedidos, y la tabla contra la que necesitas hacer join tiene 900 millones de filas. Un join simple lee la tabla grande.

La solución es calcular la lista de claves primero y empujarla hacia la lectura como un literal:

[
{ "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", "last_event_at", "channel"] }
],
"on": "customer_id",
"how": "left"
}
]

{{pending_customers}} se sustituye dentro del filtro que se ejecuta en el lado derecho, de modo que Delta y Parquet pueden usar sus estadísticas para omitir archivos en lugar de escanear la tabla. Esta es la forma declarativa del patrón collect() + isin() que usan los jobs de Spark escritos a mano.

  • El formateo es automático: las cadenas se vuelven 'a', 'b', los números 1, 2.
  • Una colección vacía se renderiza como NULL, así que IN (NULL) no coincide con nada — la respuesta correcta cuando no hay nada que enriquecer.
  • El almacén de variables se comparte con los with_transformations anidados, que es la razón por la que el lado derecho del join puede verlo.
  • Se limpia al inicio de cada ejecución, así que nada se filtra entre pipelines.

La lista de claves viaja al driver y a una cadena SQL. Decenas de miles de claves están bien; millones no — a ese tamaño, un join simple con una pista de broadcast o una tabla con buckets es la herramienta correcta. Si el conjunto de trabajo no tiene cota, omite el pushdown.

{ "type": "union", "with": { "format": "parquet", "path": "/data/orders_archive" }, "allow_missing_columns": true }

union no tiene with_transformations — el lado derecho se lee tal cual. Dale forma aguas arriba, o usa un join si lo necesitas.

Coloca un nodo debug después del join mientras desarrollas:

{ "type": "debug", "label": "after enrichment", "actions": ["count", "print_schema"] }

Un conteo de filas que creció significa que el lado derecho no era único sobre la clave del join. Ese es el bug de join más común, y una validación unique sobre la clave lo detecta de forma permanente.