Transformaciones
Las transformaciones son las entradas del array transformations. Se ejecutan en el orden en que las escribes, cada una recibiendo el DataFrame que devolvió la anterior.
Todas aceptan la meta clave skip_if_false, que activa o desactiva el paso por ejecución.
Filtrado y modelado
Sección titulada «Filtrado y modelado»Conserva las filas que cumplen una expresión booleana SQL.
{ "type": "filter", "condition": "status = 'CONFIRMED' AND amount > 0" }| Clave | Tipo | Obligatorio |
|---|---|---|
condition |
expresión SQL | sí |
Portador habitual de placeholders {{runtime}}: "id IN ({{ids}})".
Proyecta una lista de columnas — o expresiones SQL completas con alias.
{ "type": "select", "columns": ["id", "customer", "to_json(payload) AS value"] }Elimina columnas. Los nombres que no existen los ignora Spark en silencio.
{ "type": "drop", "columns": ["tmp_flag", "debug_note"] }Renombra columnas mediante un mapa ordenado, aplicado de a una.
{ "type": "rename", "mappings": { "created_at": "creation_date", "nm": "name" } }Como es secuencial, los renombramientos encadenados funcionan (a→b y luego b→c termina como c) y reordenar las claves cambia el resultado.
Convierte columnas a tipos de Spark.
{ "type": "cast", "columns": { "amount": "decimal(18,2)", "ordered_at": "date" } }Usa F.col, así que la columna debe existir. Los valores que no se pueden convertir pasan a null — Spark no lanza error.
Ordena el DataFrame.
{ "type": "sort", "columns": ["ordered_at", "id"], "ascending": true }ascending acepta un booleano o una lista de booleanos, uno por columna.
drop_duplicates
Sección titulada «drop_duplicates»Elimina filas duplicadas, opcionalmente acotado a un subconjunto.
{ "type": "drop_duplicates", "columns": ["id"] }columns es opcional: omitirlo (o pasar una lista vacía) deduplica sobre todas las columnas, exactamente como distinct. Qué fila sobrevive no es determinista — ordena antes con sort si importa.
distinct
Sección titulada «distinct»Elimina filas duplicadas usando todas las columnas. Sin parámetros.
{ "type": "distinct" }fill_na
Sección titulada «fill_na»Reemplaza valores nulos.
{ "type": "fill_na", "value": 0, "columns": ["quantity"] }{ "type": "fill_na", "value": { "quantity": 0, "segment": "UNKNOWN" } }La forma escalar toma un subconjunto opcional columns. Spark ignora los rellenos cuyo tipo no coincide con el tipo de la columna — un no-op silencioso, no un error.
Cálculo de columnas
Sección titulada «Cálculo de columnas»with_column
Sección titulada «with_column»Agrega o reemplaza columnas calculadas a partir de expresiones SQL. Dos formas mutuamente excluyentes.
// single column{ "type": "with_column", "column": "revenue", "expression": "quantity * unit_price" }
// several, in order — later expressions can use earlier ones{ "type": "with_column", "columns": { "revenue": "quantity * unit_price", "revenue_brl": "revenue * fx_rate" } }name se acepta como alias heredado de column, parseado por compatibilidad hacia atrás; escribe column en los pipelines nuevos.
Construye una columna struct anidada a partir de un mapa de campos. Más legible que un named_struct escrito a mano.
{ "type": "struct", "column": "payload", "fields": { "external_id": "contract_id", "issuer.name": "issuer_name", "issuer.document": "lpad(cast(document as string), 14, '0')", "amounts": { "principal": "principal_amount", "interest": "interest_amount" } }}- Los valores string son expresiones SQL; los valores objeto anidan más.
- Los dot-paths auto-anidan:
issuer.nameeissuer.documentse vuelven un solo structissuer, así que el payload se lee como una tabla plana en el archivo y llega anidado en los datos. - El orden de los campos sigue el orden de las claves, lo que mantiene estables los diffs.
- Un nombre de campo que contenga un punto literal es imposible — cada punto significa anidamiento.
Dos conflictos se lanzan en tiempo de aplicación: usar un segmento de ruta que ya es una hoja (a usado tanto como valor como prefijo), y duplicar una clave hoja.
Ejecuta Spark SQL arbitrario sobre el DataFrame actual.
{ "type": "sql", "view_name": "_df", "query": "SELECT customer, SUM(revenue) AS total FROM _df GROUP BY customer"}El DataFrame se registra como una vista temporal llamada view_name (por defecto _df) y el resultado de la query se vuelve el nuevo DataFrame. La vía de escape para todo lo que las demás transformaciones no expresan.
Agregación
Sección titulada «Agregación»group_by
Sección titulada «group_by»Agrupa y agrega, con pivot opcional.
{ "type": "group_by", "by": ["customer_id", "country"], "agg": [ "sum(revenue) as revenue_total", "count(*) as orders", "max(ordered_at) as last_order" ], "pivot": { "column": "month", "values": ["jan", "feb", "mar"] }}| Clave | Tipo | Obligatorio |
|---|---|---|
by |
lista de columnas | sí |
agg |
lista de expresiones SQL de agregación completas, con alias | sí |
pivot |
nombre de columna, o { column, values } |
no |
Combinación de orígenes
Sección titulada «Combinación de orígenes»Une un segundo origen leído en línea.
{ "type": "join", "with": { "format": "delta", "path": "sales.customers" }, "with_transformations": [ { "type": "filter", "condition": "active = true" }, { "type": "select", "columns": ["customer_id", "segment"] } ], "on": "customer_id", "how": "left"}| Clave | Tipo | Notas |
|---|---|---|
with |
config de origen | cualquier formato legible |
on |
string, lista o expresión SQL | "id", ["a","b"], o "l.id = r.id AND l.dt = r.dt" |
how |
tipo de join | inner (por defecto), left, right, full, cross, leftsemi, leftanti, … |
broadcast |
true / "right" / "left" / false |
hint de join map-side (broadcast) |
with_transformations |
lista | aplicadas al lado derecho antes del join |
El DataFrame izquierdo lleva el alias l y el derecho r, así que un on con expresión puede desambiguar columnas. Dentro de with_transformations los alias aún no existen — usa allí nombres de columna sin prefijo.
Broadcast join
Sección titulada «Broadcast join»Define broadcast cuando un lado es lo bastante pequeño para caber en la memoria de cada executor (una dimensión, un lookup). Spark envía ese lado a cada executor y hace el join map-side, saltándose por completo el shuffle del lado grande.
{ "type": "join", "with": { "format": "delta", "path": "ref.dim_product" }, "on": "product_id", "how": "left", "broadcast": true}| Valor | Hace broadcast de |
|---|---|
true o "right" |
el segundo origen (el lado with — la dimensión/lookup pequeña) |
"left" |
el DataFrame principal |
false o ausente |
sin hint — Spark decide por tamaño |
Hacer broadcast de un lado que en realidad no es pequeño puede agotar la memoria del executor; déjalo ausente en caso de duda.
En Studio, with y with_transformations provienen de la segunda entrada del nodo, no de un campo del formulario.
Anexa las filas de otro origen.
{ "type": "union", "with": { "format": "parquet", "path": "/data/orders_archive" }, "allow_missing_columns": false}union no tiene with_transformations: el lado derecho se lee tal cual.
Layout y número de archivos
Sección titulada «Layout y número de archivos»repartition
Sección titulada «repartition»Redistribuye las particiones del DataFrame — cambia el costo, nunca los datos.
{ "type": "repartition", "num_partitions": 64, "columns": ["pmod(hash(id), 64)"] }| Clave | Valores | Por defecto |
|---|---|---|
num_partitions |
número de particiones objetivo (entero positivo) | — |
columns |
nombres de columna o expresiones SQL; los valores iguales caen en la misma partición | — |
coalesce |
fusiona sin shuffle — solo reduce | false |
range |
repartitionByRange: divide por rango de valor en vez de por hash |
false |
Se requiere al menos uno entre num_partitions y columns.
Es la pieza que faltaba entre el partition_by de la salida, que decide qué directorios existen, y el número de archivos escritos, que nadie estaba decidiendo. Se escribe un archivo por par (task, directorio) que contiene filas, así que 200 particiones de shuffle sobre 30 días de dt dejan hasta 6.000 archivos. Reparticionar por las mismas expresiones del partition_by del destino lo reduce a un archivo por valor de clave, porque un valor de clave nunca se divide entre tasks — AQE fusiona particiones vecinas, pero nunca separa una.
[ { "type": "repartition", "columns": ["dt"] }, { "type": "repartition", "num_partitions": 1, "coalesce": true }, { "type": "repartition", "num_partitions": 8, "columns": ["fecha_evento"], "range": true }]Cada combinación inválida lanza un error con el motivo en vez de hacer otra cosa en silencio: coalesce con columns (no hay clave para agrupar), coalesce con range, coalesce sin conteo, range sin columnas, un conteo no entero o no positivo, y ningún parámetro.
Control e inspección
Sección titulada «Control e inspección»checkpoint
Sección titulada «checkpoint»Materializa el DataFrame y trunca su plan lógico.
{ "type": "checkpoint", "method": "localCheckpoint", "eager": true }| Clave | Valores | Por defecto |
|---|---|---|
method |
localCheckpoint, checkpoint |
localCheckpoint |
eager |
booleano | true |
Úsalo después de joins pesados, antes de un collect, y antes de abanicar hacia varios destinos — impide que Spark recompute el mismo linaje una y otra vez. Un method inválido se ignora, con una advertencia emitida al final de la ejecución.
collect
Sección titulada «collect»Recolecta los valores distintos de una columna en una variable de runtime.
{ "type": "collect", "column": "customer_id", "as": "active_customers" }El DataFrame pasa sin cambios, pero los valores aterrizan en {{active_customers}} para pasos posteriores. Dispara una acción del lado del driver, así que ejecútalo después de un checkpoint. Consulta variables de runtime.
| Clave | Valores | Por defecto |
|---|---|---|
column |
la columna cuyos valores distintos se recolectan | — |
as |
nombre de la variable de runtime | — |
max_values |
techo de valores distintos; 0 lo desactiva |
10000 |
La lista recolectada se convierte en un literal dentro de IN (...), y a partir de algunos miles de valores el remedio pasa a ser el problema: el plan crece, Catalyst gasta su tiempo analizando el predicado y el pushdown se degrada. Por encima de max_values el paso falla y nombra la alternativa — un join semi/inner contra la lista como DataFrame, que Spark resuelve como broadcast join sin que el driver retenga los valores. El techo se aplica dentro de la consulta (limit(max_values + 1)), así que una columna con millones de valores distintos nunca se materializa en el driver, y por encima de 1.000 valores todavía funciona, pero avisa.
stop_if_empty
Sección titulada «stop_if_empty»Termina la ejecución de forma elegante cuando no hay nada que procesar.
{ "type": "stop_if_empty", "message": "No approved orders in the window" }Las transformaciones restantes y toda escritura se saltan. El resultado regresa con skipped: true, success: true y rows_written: 0 — un no-op, no una falla. Colócalo justo después del filtro que define el conjunto de trabajo, antes de joins costosos.
Inspecciona el DataFrame sin modificarlo.
{ "type": "debug", "label": "after enrichment", "actions": ["count", "print_schema", "show"], "transformations": [{ "type": "filter", "condition": "id = 'X1'" }], "show_rows": 20, "truncate": true, "vertical": false, "extended": false}actions acepta count, print_schema, show, explain, pushdown, columns, dtypes. Sus transformations anidadas se aplican a una copia descartable usada solo para la inspección — el DataFrame del pipeline siempre pasa intacto.
pushdown
Sección titulada «pushdown»Responde lo que explain solo insinúa: ¿qué llegó realmente a la fuente?
{ "type": "debug", "label": "lectura", "actions": ["pushdown"] }Lee el plan físico e informa, por nodo de lectura, PartitionFilters (particiones podadas antes de abrir un archivo), PushedFilters (el predicado entregado a Parquet/ORC o a la base de datos), PushedAggregates, PushedGroupBy, RuntimeFilters (dynamic partition pruning y bloom filter de join) y cuántas columnas devuelve el scan. Un scan que no empujó nada se señala junto con qué hacer al respecto, y los nodos Filter por encima de los scans se cuentan — un predicado evaluado después de la lectura es dato que salió del disco para ser descartado. No dispara ningún job: el plan físico es planificación, no ejecución.
$include
Sección titulada «$include»No es una transformación sino una directiva: incrusta en línea un fragmento JSON.
{ "$include": "shared/standard_filters.json" }La ruta es relativa al archivo del pipeline. El fragmento puede ser un objeto único o una lista. Los includes anidados no se expanden, y la directiva funciona solo en el array transformations de nivel superior.
Índice rápido
Sección titulada «Índice rápido»| Tipo | Propósito |
|---|---|
filter |
conservar las filas que coinciden |
select |
proyectar columnas o expresiones |
drop |
eliminar columnas |
rename |
renombrar columnas |
cast |
cambiar tipos de columnas |
with_column |
calcular columnas |
struct |
construir una columna anidada |
drop_duplicates |
deduplicar, opcionalmente por subconjunto |
distinct |
deduplicar sobre todas las columnas |
sort |
ordenar filas |
fill_na |
reemplazar nulos |
sql |
Spark SQL arbitrario |
group_by |
agregar, con pivot opcional |
join |
unir un segundo origen |
union |
anexar otro origen |
repartition |
redistribuir particiones, controlar el número de archivos |
checkpoint |
materializar y truncar el plan |
collect |
publicar una variable de runtime |
stop_if_empty |
terminar la ejecución cuando no hay datos |
debug |
inspeccionar sin modificar |