Rendimiento
La mayoría de los pipelines lentos lo son por una de cuatro razones: leen de más, recomputan el mismo trabajo, traen datos al driver o hacen fan-out sin materializar. Cada una tiene una solución directa en el lenguaje.
Lee menos
Sección titulada «Lee menos»La fila más barata es la que nunca sale de la fuente.
// tabla de lake particionada: filtra por la columna de partición primero{ "type": "filter", "condition": "ordered_at >= '{since}'" }
// luego reduce las columnas, antes de cualquier join{ "type": "select", "columns": ["id", "customer_id", "amount"] }El orden importa: un select después del join sigue arrastrando todas las columnas por el shuffle.
Hazlo un hábito: empieza cada cadena con filter y luego select, antes de cualquier join, struct, group_by o window. El optimizador de Spark empuja parte de esto por su cuenta, pero hacerlo explícito reduce los datos para cada paso siguiente, ayuda al planner y hace que el pipeline se lea como “angosta temprano, trabaja después”.
Reutiliza la entrada sin releer
Sección titulada «Reutiliza la entrada sin releer»Cuando necesitas hacer self-join de la entrada — o correr SQL sobre ella — sin leer la fuente dos veces, regístrala una vez como temp view cacheada con input_view:
fw.run("orders.json", input_view="orders") # temp view de sesión, cacheadafw.run("orders.json", input_view={"name": "orders", "type": "global"}) # global_temp.orders// un paso posterior une la entrada consigo misma por la view — sin segunda lectura{ "type": "join", "with": { "format": "view", "path": "orders" }, "with_transformations": [ { "type": "group_by", "by": ["customer_id"], "agg": ["max(amount) as top"] } ], "on": "customer_id", "how": "left" }El DataFrame de entrada se cache()a y se registra antes de correr las transformaciones, así que el self-join reutiliza las filas cacheadas en vez de recomputar el linaje de la fuente. Defínelo en el constructor Sparquet(...) para que aplique a cada ejecución, o por llamada de run(...).
Haz checkpoint del plan
Sección titulada «Haz checkpoint del plan»{ "type": "checkpoint", "method": "localCheckpoint", "eager": true }Tres lugares lo justifican:
- Después de joins pesados, para que el plan no siga creciendo.
- Antes de
collect, una action en el driver que si no recomputaría todo el linaje. - Antes de repartir a varios destinos, para que el trabajo ocurra una vez en lugar de una por destino.
Empuja la lista de claves
Sección titulada «Empuja la lista de claves»Cuando el conjunto de trabajo es pequeño y la tabla que enriqueces es enorme:
[ { "type": "checkpoint" }, { "type": "collect", "column": "customer_id", "as": "customers" }, { "type": "join", "with": { "format": "delta", "path": "sales.bronze_events" }, "with_transformations": [ { "type": "filter", "condition": "customer_id IN ({{customers}})" } ], "on": "customer_id" }]El IN (...) literal deja que Delta y Parquet salten archivos usando sus estadísticas. Mantén la lista en decenas de miles de claves como máximo — más allá de eso, un join normal es la herramienta correcta. Ver Joins y pushdown.
Prueba que el pushdown ocurrió
Sección titulada «Prueba que el pushdown ocurrió»{ "type": "debug", "label": "lectura", "actions": ["pushdown"] }explain lo insinúa; esto lo responde. Por nodo de lectura imprime PartitionFilters, PushedFilters, PushedAggregates, PushedGroupBy, RuntimeFilters y cuántas columnas devuelve el scan, y cuenta los nodos Filter por encima de los scans — un predicado evaluado después de la lectura es dato que salió del disco para ser descartado. Lo que deja a la vista:
- un filtro sobre una columna de partición aterriza en
PartitionFiltersy salta directorios completos; un filtro sobre una columna común aterriza enPushedFiltersy solo salta los row groups que las estadísticas min/max descartan. - un predicado que la fuente no sabe expresar — una UDF, un cast,
lower(col) = 'x'— se queda comoFilterencima del scan. Reescríbelo con la columna desnuda en uno de los lados. - en JDBC, solo la lectura de tabla empuja. Con
query, el SELECT que escribiste ya es el corte. RuntimeFilterses dynamic partition pruning: la tabla de hechos podada por el filtro de la dimensión, decidido en runtime.
No cuesta nada — leer el plan físico es planificación, no ejecución.
Palancas de archivo
Sección titulada «Palancas de archivo»Los valores por defecto sirven para la mayoría de los pipelines. Estas son las que vale conocer cuando no sirven; defínelas en spark.configs.
| Config | Por defecto | Qué hace |
|---|---|---|
spark.sql.parquet.filterPushdown |
true |
empuja el predicado a las estadísticas de row group de Parquet |
spark.sql.parquet.aggregatePushdown |
false |
responde min/max/count desde el footer, sin leer datos |
spark.sql.orc.filterPushdown / .aggregatePushdown |
true / false |
lo mismo, para ORC |
spark.sql.optimizer.dynamicPartitionPruning.enabled |
true |
poda la tabla de hechos con el filtro de la dimensión en runtime |
spark.sql.optimizer.runtime.bloomFilter.enabled |
true |
bloom filter en la clave del join, recorta lo que va al shuffle |
spark.sql.files.maxPartitionBytes |
128MB |
tamaño de la task al leer archivos |
spark.sql.adaptive.coalescePartitions.enabled |
true |
fusiona particiones diminutas tras el shuffle |
aggregatePushdown es la que vale probar a mano: con ella activada, un count(1) o un max(dt) sobre Parquet se responde desde los footers de los archivos y no lee dato alguno — pero viene desactivada porque solo se aplica a una agregación desnuda sobre el scan, y cualquier filtro o cast en medio devuelve la ganancia en silencio.
Opciones de lectura de la misma familia: mergeSchema (desactivado por defecto, y cada archivo extra leído cuesta un listado), basePath (de qué prefijo se derivan las columnas de partición), recursiveFileLookup y pathGlobFilter (leer menos por no listar), y spark.sql.sources.partitionOverwriteMode: dynamic, que hace que overwrite reemplace solo las particiones presentes en el DataFrame en vez del directorio completo.
Detente temprano
Sección titulada «Detente temprano»{ "type": "stop_if_empty", "message": "Nada que procesar" }Colocado justo después del filtro que define el conjunto de trabajo, convierte una noche de no-op en segundos en lugar de un pipeline entero sobre cero filas.
Modela el lado derecho de un join
Sección titulada «Modela el lado derecho de un join»"with_transformations": [ { "type": "filter", "condition": "active = true" }, { "type": "select", "columns": ["customer_id", "segment"] }, { "type": "distinct" }]Filtra, proyecta, deduplica — en ese orden — antes del join. Un lado derecho que repite la clave multiplica filas, lo que es a la vez un bug de correctitud y de rendimiento.
El skew son dos problemas distintos
Sección titulada «El skew son dos problemas distintos»- Skew de join, en la lectura — una clave con muchas más filas que el resto deja una sola task corriendo durante minutos mientras las demás ya terminaron. Este lo resuelve AQE:
spark.sql.adaptive.skewJoin.enabled(por defectotrue), conskewedPartitionFactor(5) yskewedPartitionThresholdInBytes(256MB) definiendo qué cuenta como skew; divide la partición problemática y replica el lado correspondiente. Solo actúa en un sort-merge join con AQE activado, así que un join con hint de broadcast nunca llega ahí. - Skew de escritura — un valor de la columna de
partition_byque concentra la mayoría de las filas escribe un archivo enorme mientras los demás directorios ya terminaron. Aquí AQE no ayuda, porque el layout de directorios es el requisito, no un accidente del plan. Divídelo a propósito:repartitionpor las columnas de partición más una sal (pmod(hash(id), 8)), lo que convierte el archivo único de ese directorio en ocho — y el lector no paga nada por los archivos extra mientras sigan siendo grandes.
Cuida la escritura
Sección titulada «Cuida la escritura»partition_bycon alta cardinalidad crea miles de archivos diminutos. Particiona por algo grueso (fecha, país), nunca por un id.mergecuesta más queappend. Úsalo cuando necesites idempotencia, no por defecto.inferSchemade CSV cuesta una pasada extra. Definecastexplícito y desactívalo para feeds grandes y estables.
DataFusion Comet, opt-in
Sección titulada «DataFusion Comet, opt-in»Comet es un plugin de Spark que reemplaza operadores del plan físico por implementaciones nativas (Rust/Arrow) y vuelve a Spark para lo que no soporta. El framework no necesita código para ello — es configuración de sesión:
Por qué activarlo. El beneficio no es “Rust es rápido” — son cuatro propiedades concretas:
- Ejecución vectorizada sobre Arrow. El operador nativo trabaja en lotes columnares, con instrucciones SIMD, en lugar de fila por fila en la JVM. Por eso la ganancia aparece en el escaneo, el filtro, la proyección, el hash aggregate y el sort, y desaparece donde el tiempo es de red.
- Memoria fuera del heap. Los buffers de ejecución salen del heap de la JVM, lo que le quita al recolector de basura justamente lo que más asigna. En un job con pausas de GC visibles en la timeline de Spark, esa es la mitad de la ganancia que se nota primero.
- Ningún cambio en el pipeline. Es configuración de sesión, no API: el mismo JSON, los mismos readers, writers, transformaciones y validaciones. Nada en el archivo del pipeline dice que Comet existe.
- Adopción parcial y reversible. El reemplazo es operador por operador, y lo que Comet no soporta sigue en Spark, dentro del mismo plan. No hay migración ni puerta de un solo sentido: quitar las configuraciones vuelve al estado anterior en la siguiente ejecución.
Compensa en lecturas voluminosas de Parquet con filtro, proyección, agregación, sort o shuffle en el camino — el perfil de la mayor parte de la ingesta por lotes. No compensa cuando el tiempo está del otro lado de la red (JDBC, APIs), con UDF en Python (que fuerza ese tramo de vuelta a Spark), en pipelines dominados por la escritura, ni en volúmenes pequeños, donde dimensionar el off-heap y cargar 88 MB de jar cuesta más de lo que ahorra.
{ "spark": { "configs": { "spark.plugins": "org.apache.spark.CometPlugin", "spark.shuffle.manager": "org.apache.spark.sql.comet.execution.shuffle.CometShuffleManager", "spark.memory.offHeap.enabled": "true", "spark.memory.offHeap.size": "4g" } }}Medido en el pipeline del propio framework, 40.000.000 de filas / 290 MB de Parquet, local[4], tres repeticiones cronometradas tras un calentamiento descartado (mediana):
| Forma | Sin Comet | Con Comet | Ganancia |
|---|---|---|---|
filtro + group_by con sum/count |
3,56s | 1,61s | 2,21x |
filtro + count |
1,02s | 0,83s | 1,24x |
La ganancia sigue cuánto del plan pasó a nativo: la agregación volvió con CometNativeScan, CometFilter, CometHashAggregate, CometExchange y CometNativeShuffle, mientras que la forma de conteo tiene mucho menos que acelerar. Antes de activarlo:
- el jar tiene que estar en el classpath del driver antes de que arranque la JVM —
--driver-class-path,spark.driver.extraClassPatho el directorio de jars del clúster. Porspark.jarsllega tarde y la sesión muere conjava.lang.ClassNotFoundException: org.apache.spark.CometPlugin. - el artefacto es por línea de Spark:
org.apache.datafusion:comet-spark-spark4.1_2.13:1.0.0para Spark 4.1, unos 88 MB. - la biblioteca nativa existe solo para Linux. En cualquier otro sistema el plugin carga, se desactiva en silencio y el pipeline pasa sin aceleración alguna — así que revisa los nodos
Comet*en el plan en vez de suponerlo. - la memoria off-heap es obligatoria (
spark.memory.offHeap.enabledy.size). - el fallback es por operador, y
spark.comet.explain.fallback.enabledregistra el motivo de cada uno.
Mide antes de adivinar
Sección titulada «Mide antes de adivinar»{ "type": "debug", "label": "after join", "actions": ["count", "explain"], "extended": false }explain muestra el plan que Spark realmente ejecutará — la forma más rápida de confirmar que un filtro se empujó o que ocurrió un broadcast. count tras cada etapa dice dónde se multiplicaron las filas.
El panel Run en Studio reporta filas leídas, filas escritas y duración por ejecución, lo que suele bastar para detectar la regresión sin abrir la Spark UI.
Una lista rápida
Sección titulada «Una lista rápida»| Síntoma | Causa probable |
|---|---|
| El tiempo escala con el número de salidas | Falta checkpoint antes del fan-out |
| Driver sin memoria | collect en una columna grande, o una lista IN (...) enorme |
| El join tarda más que todo lo demás | Lado derecho no filtrado/proyectado antes del join |
| Miles de archivos diminutos | partition_by en una columna de alta cardinalidad |
| Lento cada noche, aun sin datos | Falta stop_if_empty |
Un Filter aparece encima del scan en el plan |
Predicado que la fuente no expresa — confírmalo con debug / pushdown |
| Una task corre durante minutos mientras las demás terminaron | Skew de join (AQE), o un valor de partición que concentra las filas (escritura) |