Ir al contenido
SparquetSparquet

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.

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”.

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, cacheada
fw.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(...).

{ "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.

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.

{ "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 PartitionFilters y salta directorios completos; un filtro sobre una columna común aterriza en PushedFilters y 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 como Filter encima 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.
  • RuntimeFilters es 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.

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.

{ "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.

"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.

  • 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 defecto true), con skewedPartitionFactor (5) y skewedPartitionThresholdInBytes (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_by que 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: repartition por 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.
  • partition_by con alta cardinalidad crea miles de archivos diminutos. Particiona por algo grueso (fecha, país), nunca por un id.
  • merge cuesta más que append. Úsalo cuando necesites idempotencia, no por defecto.
  • inferSchema de CSV cuesta una pasada extra. Define cast explícito y desactívalo para feeds grandes y estables.

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.extraClassPath o el directorio de jars del clúster. Por spark.jars llega tarde y la sesión muere con java.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.0 para 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.enabled y .size).
  • el fallback es por operador, y spark.comet.explain.fallback.enabled registra el motivo de cada uno.
{ "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.

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)