Saltar al contenido

lección 5

JOINs distribuidos: shuffle, broadcast y por qué importa

El shuffle es el villano del rendimiento. El broadcast join es tu arma secreta cuando una tabla es pequeña.

60 min

### El JOIN: fácil en SQL, complicado en un cluster

En SQL tradicional (PostgreSQL en tu portátil), un JOIN es simple: la base de datos tiene todas las tablas en el mismo disco, lee ambas, las cruza en memoria, y te da el resultado. Pero en un sistema distribuido como Spark, las dos tablas están repartidas entre múltiples nodos. La tabla de ventas tiene particiones en el nodo 1, 2, 3... y la tabla de clientes tiene particiones en el nodo 4, 5, 6... Para cruzarlas, los datos tienen que MOVERSE entre nodos.

Este movimiento de datos entre nodos se llama SHUFFLE y es, sin exagerar, la operación más cara en procesamiento distribuido. Es el momento en que la red se convierte en el cuello de botella. Entender cuándo ocurre un shuffle, por qué es caro, y cómo evitarlo es la diferencia entre un job que tarda 5 minutos y uno que tarda 3 horas.

### El shuffle: el villano del rendimiento

Imagina que tienes 4 amigos, cada uno con una caja de cartas de Pokémon desordenadas. Quieres organizar TODAS las cartas por tipo (fuego, agua, planta...). Cada amigo tendría que pasar sus cartas de fuego al amigo encargado de fuego, las de agua al de agua, etc. Todos pasan cartas a todos — es un caos logístico. Eso es un shuffle.

En Spark, un shuffle significa: cada nodo envía una porción de sus datos a TODOS los demás nodos del cluster, reorganizando los datos por la clave del JOIN o del groupBy. Si tienes 10 nodos con 10 GB cada uno, un shuffle puede mover 100 GB de datos por la red. La red es mucho más lenta que la RAM — incluso con 10 Gbps, mover 100 GB tarda 80 segundos solo en transferencia, sin contar la serialización, deserialización, y escritura temporal a disco.

El shuffle reorganiza datos entre nodos por clave. Es necesario para JOINs y groupBy, pero es MUY caro.

### Sort-Merge Join: el join por defecto

Cuando haces un JOIN entre dos DataFrames grandes, Spark usa por defecto el Sort-Merge Join. El proceso es: 1) Shuffle: reorganiza ambas tablas por la clave del JOIN (los datos de la misma clave van al mismo nodo). 2) Sort: ordena cada partición por la clave. 3) Merge: recorre ambas particiones ordenadas en paralelo, emparejando filas.

El Sort-Merge Join es robusto y funciona con tablas de cualquier tamaño. Pero el shuffle inicial es caro. Si tu tabla de ventas tiene 100 millones de filas y tu tabla de clientes tiene 5 millones, AMBAS se redistribuyen por la red. Mueves 105 millones de filas entre nodos.

1from pyspark.sql import SparkSession
2from pyspark.sql.functions import col
3
4spark = SparkSession.builder.appName("joins").master("spark://spark-master:7077").getOrCreate()
5
6# Tabla grande: 1 millón de ventas
7ventas = spark.read.option("header", "true").option("inferSchema", "true").csv("/home/jovyan/data/ventas.csv")
8
9# Tabla mediana: 50.000 clientes
10clientes = spark.createDataFrame(
11 [(i, f"Cliente_{i}", f"ciudad_{i % 10}") for i in range(1, 50001)],
12 ["id", "nombre", "ciudad_cliente"]
13)
14
15# JOIN normal → Sort-Merge Join (shuffle de AMBAS tablas)
16resultado = ventas.join(clientes, ventas.cliente_id == clientes.id, "inner")
17
18resultado.explain()
19# == Physical Plan ==
20# SortMergeJoin [cliente_id], [id]
21# :- Sort [cliente_id ASC]
22# : +- Exchange hashpartitioning(cliente_id, 200) ← SHUFFLE ventas
23# : +- FileScan csv [...]
24# +- Sort [id ASC]
25# +- Exchange hashpartitioning(id, 200) ← SHUFFLE clientes
26# +- ...

Sort-Merge Join: ambas tablas sufren un shuffle (Exchange) costoso

### Broadcast Join: el arma secreta para tablas pequeñas

Aquí viene el truco que te ahorrará horas de ejecución: si una de las dos tablas es lo suficientemente pequeña para caber en la memoria de un executor (digamos, menos de 100 MB), puedes decirle a Spark que la ENVÍE COMPLETA a cada nodo del cluster. Así, cada nodo tiene la tabla pequeña entera en memoria y puede hacer el JOIN localmente, SIN mover la tabla grande. Cero shuffle en la tabla grande.

Es como la diferencia entre llevar una enciclopedia a cada biblioteca del país (broadcast) vs. llevar todos los libros de todas las bibliotecas a un punto central (shuffle). Si la enciclopedia cabe en una mochila, es mucho más eficiente repartirla que mover millones de libros.

1from pyspark.sql.functions import broadcast
2
3# BROADCAST JOIN: envía la tabla pequeña (clientes) a TODOS los nodos
4# La tabla grande (ventas) NO se mueve — cada nodo cruza su partición local
5resultado_broadcast = ventas.join(
6 broadcast(clientes), # ← la magia: envía clientes a todos los nodos
7 ventas.cliente_id == clientes.id,
8 "inner"
9)
10
11resultado_broadcast.explain()
12# == Physical Plan ==
13# BroadcastHashJoin [cliente_id], [id]
14# :- FileScan csv [...] ← ventas NO se mueve
15# +- BroadcastExchange HashedRelation(id) ← clientes se copia a cada nodo
16# +- ...
17#
18# ¡NO hay Exchange (shuffle) en ventas! Solo se mueven los 50K clientes.

Broadcast Join: la tabla pequeña se copia a cada nodo, la grande NO se mueve

Broadcast Join: la tabla pequeña viaja a todos los nodos; la tabla grande se queda quieta

### ¿Cuándo usar broadcast?

La regla es simple: si una de las dos tablas cabe en la memoria de un executor, usa broadcast. Spark lo hace automáticamente si la tabla es menor que spark.sql.autoBroadcastJoinThreshold (por defecto 10 MB). Pero puedes forzarlo con broadcast() para tablas más grandes que conozcas que son "pequeñas" en tu contexto.

  • Tabla de productos (10.000 filas, 1 MB): SIEMPRE broadcast.
  • Tabla de clientes (100.000 filas, 20 MB): broadcast si el executor tiene suficiente RAM.
  • Tabla de ventas (100 millones de filas, 50 GB): NUNCA broadcast. Es la tabla que se queda quieta.
  • Tabla de dimensiones (países, categorías, estados): broadcast siempre, son pequeñísimas.

Consejo de senior: en un modelo de data warehouse (estrella), las tablas de hechos son enormes y las de dimensiones son pequeñas. El patrón perfecto para broadcast joins: broadcast TODAS las dimensiones y deja la tabla de hechos sin mover. He visto jobs pasar de 45 minutos a 3 minutos con este cambio.

Si intentas hacer broadcast de una tabla demasiado grande, el executor se quedará sin memoria y el job fallará con OutOfMemoryError. En producción, siempre verifica el tamaño real de la tabla antes de forzar un broadcast. Usa df.persist() + spark.catalog.cacheTable() para medir.

### Comparativa de rendimiento: shuffle vs broadcast

1import time
2
3# Medir tiempo del shuffle join
4start = time.time()
5resultado_shuffle = ventas.join(clientes, ventas.cliente_id == clientes.id).count()
6tiempo_shuffle = time.time() - start
7print(f"Shuffle Join: {tiempo_shuffle:.2f}s")
8
9# Medir tiempo del broadcast join
10start = time.time()
11resultado_broadcast = ventas.join(broadcast(clientes), ventas.cliente_id == clientes.id).count()
12tiempo_broadcast = time.time() - start
13print(f"Broadcast Join: {tiempo_broadcast:.2f}s")
14
15print(f"Broadcast es {tiempo_shuffle / tiempo_broadcast:.1f}x más rápido")
16# Resultado típico: Broadcast es 3-10x más rápido para tablas pequeñas

Comparativa real: broadcast gana por 3-10x cuando la tabla cabe en memoria

### Configurar el umbral de auto-broadcast

1# Ver el umbral actual (por defecto 10 MB)
2print(spark.conf.get("spark.sql.autoBroadcastJoinThreshold"))
3# "10485760" (10 MB en bytes)
4
5# Subir el umbral a 100 MB (si tus executors tienen suficiente RAM)
6spark.conf.set("spark.sql.autoBroadcastJoinThreshold", 100 * 1024 * 1024)
7
8# Desactivar auto-broadcast (para forzar siempre shuffle — útil para testing)
9spark.conf.set("spark.sql.autoBroadcastJoinThreshold", -1)

Controlar cuándo Spark decide automáticamente hacer broadcast

## ejercicios

[01]

Elegir el tipo de JOIN correcto

Tienes una tabla de ventas (50M filas, 20 GB) y una tabla de categorías (50 filas, 2 KB). Escribe el JOIN usando broadcast en la tabla correcta y explica por qué.

Cargando editor...
[02]

Medir el impacto del shuffle

Compara el tiempo de ejecución de un JOIN con y sin broadcast. Mide ambos, imprime los tiempos y calcula la mejora.

Cargando editor...
[03]

JOIN con múltiples tablas de dimensiones

Haz un JOIN de la tabla de ventas con 3 tablas de dimensiones (productos, ciudades, fechas) usando broadcast en todas las dimensiones. Es el patrón clásico de warehouse.

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