Saltar al contenido

lección 6

Particiones: repartition, coalesce y controlando la distribución

Las particiones determinan el paralelismo. Aprende a controlar cuántas hay y cuándo ajustarlas.

55 min

### Las particiones: los ladrillos del paralelismo

Si los datos son un pastel, las particiones son los trozos en que lo cortas. Cada trozo (partición) se procesa de forma independiente en un executor. Si tienes 8 cores y 4 particiones, 4 cores estarán ociosos. Si tienes 8 cores y 800 particiones, cada core procesará ~100 particiones secuencialmente (con overhead de gestión). El número óptimo de particiones es uno de los tuning más importantes en Spark.

Cuando lees un archivo, Spark decide automáticamente cuántas particiones crear basándose en el tamaño (por defecto, una partición por cada 128 MB). Pero después de operaciones como filtros agresivos o joins, puedes acabar con particiones desbalanceadas: unas enormes y otras vacías. Eso es data skew y es un asesino del rendimiento.

### ¿Cuántas particiones tengo?

1from pyspark.sql import SparkSession
2from pyspark.sql.functions import col, spark_partition_id
3
4spark = SparkSession.builder.appName("particiones").master("spark://spark-master:7077").getOrCreate()
5ventas = spark.read.option("header", "true").option("inferSchema", "true").csv("/home/jovyan/data/ventas.csv")
6
7# Ver cuántas particiones tiene un DataFrame
8print(f"Particiones: {ventas.rdd.getNumPartitions()}")
9# Resultado típico para 1M filas: 2-8 particiones (depende del tamaño del archivo)
10
11# Ver cuántas filas hay en cada partición
12from pyspark.sql.functions import spark_partition_id
13
14ventas.withColumn("partition_id", spark_partition_id()).groupBy("partition_id").count().show()
15# +------------+------+
16# |partition_id| count|
17# +------------+------+
18# | 0|500234|
19# | 1|499766|
20# +------------+------+

Inspeccionar el número de particiones y su distribución de filas

### repartition(): redistribuir con shuffle

repartition(n) redistribuye los datos en exactamente n particiones. Causa un SHUFFLE completo (todos los datos se mueven entre nodos), pero el resultado son particiones equilibradas. Úsalo cuando:

  • Necesitas más paralelismo (tienes pocas particiones y muchos cores disponibles).
  • Quieres particionar por una columna para optimizar JOINs posteriores.
  • Después de un filtro que dejó muchas particiones vacías.
  • Antes de escribir a disco y quieres controlar el número de archivos de salida.
1# Aumentar particiones para más paralelismo
2ventas_200 = ventas.repartition(200)
3print(f"Antes: {ventas.rdd.getNumPartitions()} particiones")
4print(f"Después: {ventas_200.rdd.getNumPartitions()} particiones")
5
6# Particionar por columna: datos de la misma ciudad van a la misma partición
7ventas_por_ciudad = ventas.repartition("ciudad")
8# Esto es útil si vas a hacer groupBy("ciudad") después — evita un shuffle extra
9
10# Particionar por columna con número fijo
11ventas_por_ciudad_10 = ventas.repartition(10, "ciudad")
12# 10 particiones, agrupadas por ciudad

repartition: redistribuye con shuffle pero da particiones equilibradas

repartition() SIEMPRE causa un shuffle. No lo uses a la ligera. Si solo quieres REDUCIR el número de particiones (ej: de 200 a 10 antes de escribir), usa coalesce() que es mucho más eficiente.

### coalesce(): reducir particiones SIN shuffle

coalesce(n) reduce el número de particiones a n combinando particiones existentes en el mismo nodo. NO mueve datos entre nodos (sin shuffle). Es perfecto para reducir particiones antes de escribir archivos (no quieres 200 archivos Parquet cuando 10 son suficientes).

1# Reducir particiones sin shuffle (eficiente)
2ventas_reducido = ventas.repartition(200).coalesce(10)
3print(f"Particiones: {ventas_reducido.rdd.getNumPartitions()}") # 10
4
5# Caso de uso real: escribir exactamente 5 archivos Parquet
6(
7 ventas
8 .filter(col("ciudad") == "Madrid")
9 .coalesce(5) # 5 particiones = 5 archivos de salida
10 .write.mode("overwrite")
11 .parquet("/home/jovyan/output/ventas_madrid/")
12)
13# Resultado: 5 archivos Parquet en la carpeta de salida

coalesce: reduce particiones fusionando las existentes (sin mover datos entre nodos)

### La regla de oro para particiones

  • Particiones óptimas para procesamiento: 2-4x el número de cores del cluster.
  • Tamaño óptimo por partición: 100-200 MB. Ni muy grandes (OOM) ni muy pequeñas (overhead).
  • Antes de escribir: coalesce al número de archivos que quieres generar.
  • Después de un filtro agresivo: repartition si las particiones están vacías o desbalanceadas.
  • Para groupBy o JOIN: repartition por la clave si la vas a usar varias veces.

Consejo de senior: el valor por defecto de spark.sql.shuffle.partitions es 200. Para datasets pequeños (< 1 GB), bájalo a 10-20. Para datasets enormes (> 100 GB), súbelo a 500-2000. Un mal valor aquí es la causa #1 de jobs lentos que veo en producción.

### Data skew: el enemigo silencioso

Data skew ocurre cuando unas particiones tienen muchos más datos que otras. Ejemplo: si particionas por ciudad y el 60% de tus ventas son de Madrid, la partición de Madrid tendrá 6x más datos que las demás. El job tarda lo que tarda la partición MÁS lenta — así que un nodo trabaja 6x más que los demás mientras los otros esperan.

1# Detectar data skew: ver la distribución de filas por partición
2from pyspark.sql.functions import spark_partition_id, count
3
4ventas_por_ciudad = ventas.repartition("ciudad")
5
6distribucion = (
7 ventas_por_ciudad
8 .withColumn("pid", spark_partition_id())
9 .groupBy("pid")
10 .agg(count("*").alias("filas"))
11 .orderBy("filas", ascending=False)
12)
13
14distribucion.show()
15# Si ves una partición con 500K filas y otra con 10K → tienes skew
16
17# Solución: salted key (añadir ruido a la clave de partición)
18from pyspark.sql.functions import concat, lit, floor, rand
19
20ventas_salted = ventas.withColumn(
21 "ciudad_salted",
22 concat(col("ciudad"), lit("_"), floor(rand() * 10).cast("string"))
23)
24# Ahora "Madrid" se convierte en "Madrid_0", "Madrid_1", ..., "Madrid_9"
25# Se reparte en 10 particiones en vez de 1

Detectar y solucionar data skew con salted keys

## ejercicios

[01]

Ajustar particiones para un pipeline

Lee el CSV de ventas, verifica cuántas particiones tiene, reparticiona a 20 para procesamiento, haz un filtro y agrupación, y antes de escribir reduce a 5 archivos con coalesce.

Cargando editor...
[02]

Detectar data skew en tus datos

Particiona las ventas por "producto" y analiza la distribución. ¿Hay algún producto con muchas más ventas que otro? Calcula el ratio entre la partición más grande y la más pequeña.

Cargando editor...
[03]

Elegir entre coalesce y repartition

Dado un escenario, elige la operación correcta y explica por qué. Completa el código con coalesce o repartition según el caso.

Cargando editor...

Regístrate para guardar tu progreso.

## comentarios

Reporta erratas, ayuda a otros o comparte tu opinión. Sé constructivo.

Inicia sesión para comentar y responder.

cargando comentarios...