Ir al contenido
SparquetSparquet

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.

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

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.

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.

Elimina filas duplicadas usando todas las columnas. Sin parámetros.

{ "type": "distinct" }

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.

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.name e issuer.document se vuelven un solo struct issuer, 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.

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
agg lista de expresiones SQL de agregación completas, con alias
pivot nombre de columna, o { column, values } no

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.

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.

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.

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.

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.

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.

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.

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.

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