Módulo 5: Joins And Window Functions At Scale
Mini-proyecto: joins y rankings de Kiosko a escala
Descripción
Este proyecto cierra el módulo integrando las seis lecciones anteriores en un solo script: lees fact_orders_at_scale y las dos tablas de dimensión, confirmas con .explain() que Spark elige BroadcastHashJoin por defecto, fuerzas y confirmas SortMergeJoin desactivando el umbral, calculas el revenue acumulado por tienda con desempate explícito, y cierras con el ranking de producto top por tienda y por día. Todo, verificado con assert, sobre los diez millones de filas completos de fact_orders_at_scale.
Conexión con el módulo. Este proyecto no introduce ningún concepto nuevo — es la integración final de las lecciones 2 a 7, en el mismo orden en que las construiste, con un solo script corrido de punta a punta sobre el dataset sintético completo.
Una analogía: la auditoría completa de la temporada
Retoma las dos analogías centrales de este módulo: el directorio fotocopiado (o no) para cada JOIN, y el corredor con su tiempo acumulado y su posición en cada punto de control. Este proyecto es la auditoría completa de una temporada entera de reparto y de carrera: ¿qué decisión tomó cada camión frente a cada directorio, con evidencia de su plan de ruta? ¿Cuál fue el acumulado final de cada corredor, verificado punto por punto? ¿Quién quedó primero, en cada categoría, en cada punto de control? Una auditoría que no revisa cada una de estas preguntas, con evidencia ejecutada, no es una auditoría completa — es solo la palabra de alguien de que "todo salió bien".
El material: todo lo que este módulo construyó, en un solo flujo
Necesitas: kiosko_orders_at_scale.csv (generado en el módulo 4, lección 4, con generate_orders_at_scale(250_000)), dim_store.csv y dim_product.csv (las mismas dos tablas de dimensión del módulo 3), todos en la misma carpeta donde vas a correr este script.
La solución de referencia, verificada
Parte 1 y 2 — Abrir la sesión, leer fact_orders_at_scale y las dimensiones
# kiosko_scaled_joins_and_rankings.py
from pyspark.sql import SparkSession, Window
from pyspark.sql import functions as F
from pyspark.sql.types import (
StructType, StructField, StringType, IntegerType, DoubleType, TimestampType,
)
print("=== Kiosko a escala, joins y ventanas: entrega final del modulo 5 ===\n")
print("Parte 1 -- abriendo la SparkSession")
spark = SparkSession.builder.appName("kiosko-spark").master("local[*]").getOrCreate()
print(f"Spark version: {spark.version}\n")
print("Parte 2 (L2-L4) -- leyendo fact_orders_at_scale, dim_store, dim_product")
scale_schema = StructType([
StructField("order_id", StringType(), False),
StructField("franchise_id", IntegerType(), False),
StructField("store_id", StringType(), False),
StructField("product_id", StringType(), False),
StructField("quantity", IntegerType(), False),
StructField("unit_price", DoubleType(), False),
StructField("order_ts", TimestampType(), False),
])
dim_store_schema = StructType([
StructField("store_id", StringType(), False),
StructField("store_name", StringType(), False),
StructField("city", StringType(), False),
])
dim_product_schema = StructType([
StructField("product_id", StringType(), False),
StructField("product_name", StringType(), False),
StructField("category", StringType(), False),
StructField("unit_cost", DoubleType(), False),
])
orders_at_scale_df = spark.read.csv("kiosko_orders_at_scale.csv", schema=scale_schema, header=True, enforceSchema=False)
dim_store_df = spark.read.csv("dim_store.csv", schema=dim_store_schema, header=True, enforceSchema=False)
dim_product_df = spark.read.csv("dim_product.csv", schema=dim_product_schema, header=True, enforceSchema=False)
fact_orders_at_scale_df = orders_at_scale_df.withColumn("revenue", F.col("quantity") * F.col("unit_price"))
row_count = fact_orders_at_scale_df.count()
print(f"fact_orders_at_scale_df.count() = {row_count}")
assert row_count == 10_000_000
print("Verificacion: 10,000,000 filas -> OK\n")
Parte 3 — El JOIN por defecto: BroadcastHashJoin
print("Parte 3 (L2-L3) -- join por defecto: BroadcastHashJoin")
joined_default = fact_orders_at_scale_df.join(dim_store_df, "store_id").join(dim_product_df, "product_id")
print(f"spark.sql.autoBroadcastJoinThreshold = {spark.conf.get('spark.sql.autoBroadcastJoinThreshold')}")
print("Plan fisico (nota BroadcastHashJoin / BroadcastExchange, sin Exchange de shuffle):")
joined_default.explain()
count_default = joined_default.count()
print(f"joined_default.count() = {count_default}")
assert count_default == row_count == 10_000_000
print("Verificacion: BroadcastHashJoin por defecto, sin perder filas -> OK\n")
Parte 4 — El mismo JOIN, forzado a SortMergeJoin
print("Parte 4 (L4) -- mismo join, forzado a SortMergeJoin")
spark.conf.set("spark.sql.autoBroadcastJoinThreshold", -1)
joined_forced = fact_orders_at_scale_df.join(dim_store_df, "store_id")
print(f"spark.sql.autoBroadcastJoinThreshold = {spark.conf.get('spark.sql.autoBroadcastJoinThreshold')}")
print("Plan fisico (nota SortMergeJoin + Exchange en ambos lados):")
joined_forced.explain()
count_forced = joined_forced.count()
assert count_forced == row_count == 10_000_000
print(f"joined_forced.count() = {count_forced}")
print("Verificacion: SortMergeJoin forzado, mismo resultado, plan distinto -> OK")
spark.conf.set("spark.sql.autoBroadcastJoinThreshold", 10 * 1024 * 1024)
print(f"spark.sql.autoBroadcastJoinThreshold restaurado = {spark.conf.get('spark.sql.autoBroadcastJoinThreshold')}\n")
Parte 5 — Revenue acumulado por tienda, con desempate
print("Parte 5 (L5-L6) -- revenue acumulado por tienda, con tiebreaker")
store_window = Window.partitionBy("store_id").orderBy("order_ts", "franchise_id", "order_id")
with_running = fact_orders_at_scale_df.withColumn("running_total", F.round(F.sum("revenue").over(store_window), 2))
final_by_store = {
r["store_id"]: r["max_running"]
for r in with_running.groupBy("store_id").agg(F.max("running_total").alias("max_running")).collect()
}
print(f"running_total final por tienda = {final_by_store}")
assert final_by_store == {"S01": 9575000.0, "S02": 9700000.0, "S03": 7262500.0}
print("Verificacion: running_total final coincide con el desglose conocido -> OK\n")
Parte 6 — Producto top por tienda y por día
print("Parte 6 (L7) -- producto top por tienda y por dia")
fact_with_day = fact_orders_at_scale_df.withColumn("order_day", F.to_date("order_ts"))
revenue_by_product_day = (
fact_with_day.groupBy("store_id", "order_day", "product_id")
.agg(F.round(F.sum("revenue"), 2).alias("total_revenue"))
)
rank_window = Window.partitionBy("store_id", "order_day").orderBy(F.desc("total_revenue"), "product_id")
top_products = (
revenue_by_product_day
.withColumn("rank", F.row_number().over(rank_window))
.filter(F.col("rank") == 1)
.orderBy("store_id", "order_day")
)
num_top = top_products.count()
print(f"num_top (filas con rank == 1) = {num_top}")
assert num_top == 20
sample = top_products.filter((F.col("store_id") == "S01") & (F.col("order_day") == "2026-08-08")).collect()[0]
print(f"S01 2026-08-08 -- producto top = {sample['product_id']}, total_revenue = {sample['total_revenue']}")
assert sample["product_id"] == "P004"
assert sample["total_revenue"] == 2_250_000.0
print("Verificacion: 20 combinaciones store x dia, cada una con su producto top, revenue escalado x250,000 -> OK\n")
Parte 7 — Resumen final
print("Parte 7 -- resumen final")
total_revenue = round(sum(final_by_store.values()), 2)
print(f"Filas procesadas: {row_count:,}")
print(f"Revenue total verificado: {total_revenue:,}")
assert total_revenue == 26_537_500.00
spark.stop()
print("=== spark.stop() -- modulo 5 cerrado, joins y ventanas verificados sobre 10M filas ===")
Qué esperar. Al correr python3 kiosko_scaled_joins_and_rankings.py completo (las siete partes juntas), la salida es exactamente esta (ejecutada en esta corrida, PySpark 4.2.0):
=== Kiosko a escala, joins y ventanas: entrega final del modulo 5 ===
Parte 1 -- abriendo la SparkSession
Spark version: 4.2.0
Parte 2 (L2-L4) -- leyendo fact_orders_at_scale, dim_store, dim_product
fact_orders_at_scale_df.count() = 10000000
Verificacion: 10,000,000 filas -> OK
Parte 3 (L2-L3) -- join por defecto: BroadcastHashJoin
spark.sql.autoBroadcastJoinThreshold = 10485760b
Plan fisico (nota BroadcastHashJoin / BroadcastExchange, sin Exchange de shuffle):
== Physical Plan ==
AdaptiveSparkPlan isFinalPlan=false
+- Project [product_id#3, store_id#2, order_id#0, franchise_id#1, quantity#4, unit_price#5, order_ts#6, revenue#15, store_name#8, city#9, product_name#11, category#12, unit_cost#13]
+- BroadcastHashJoin [product_id#3], [product_id#10], Inner, BuildRight, false, false
:- Project [store_id#2, order_id#0, franchise_id#1, product_id#3, quantity#4, unit_price#5, order_ts#6, revenue#15, store_name#8, city#9]
: +- BroadcastHashJoin [store_id#2], [store_id#7], Inner, BuildRight, false, false
: :- Project [order_id#0, franchise_id#1, store_id#2, product_id#3, quantity#4, unit_price#5, order_ts#6, (cast(quantity#4 as double) * unit_price#5) AS revenue#15]
: : +- Filter (isnotnull(store_id#2) AND isnotnull(product_id#3))
: : +- FileScan csv [order_id#0,franchise_id#1,store_id#2,product_id#3,quantity#4,unit_price#5,order_ts#6] Batched: false, DataFilters: [isnotnull(store_id#2), isnotnull(product_id#3)], Format: CSV, ...
: +- BroadcastExchange HashedRelationBroadcastMode(List(input[0, string, false]),false), [plan_id=75]
: +- Filter isnotnull(store_id#7)
: +- FileScan csv [store_id#7,store_name#8,city#9] Batched: false, DataFilters: [isnotnull(store_id#7)], Format: CSV, ...
+- BroadcastExchange HashedRelationBroadcastMode(List(input[0, string, false]),false), [plan_id=79]
+- Filter isnotnull(product_id#10)
+- FileScan csv [product_id#10,product_name#11,category#12,unit_cost#13] Batched: false, DataFilters: [isnotnull(product_id#10)], Format: CSV, ...
joined_default.count() = 10000000
Verificacion: BroadcastHashJoin por defecto, sin perder filas -> OK
Parte 4 (L4) -- mismo join, forzado a SortMergeJoin
spark.sql.autoBroadcastJoinThreshold = -1
Plan fisico (nota SortMergeJoin + Exchange en ambos lados):
== Physical Plan ==
AdaptiveSparkPlan isFinalPlan=false
+- Project [store_id#2, order_id#0, franchise_id#1, product_id#3, quantity#4, unit_price#5, order_ts#6, revenue#15, store_name#8, city#9]
+- SortMergeJoin [store_id#2], [store_id#7], Inner
:- Sort [store_id#2 ASC NULLS FIRST], false, 0
: +- Exchange hashpartitioning(store_id#2, 200), ENSURE_REQUIREMENTS, [plan_id=331]
: +- Project [order_id#0, franchise_id#1, store_id#2, product_id#3, quantity#4, unit_price#5, order_ts#6, (cast(quantity#4 as double) * unit_price#5) AS revenue#15]
: +- Filter isnotnull(store_id#2)
: +- FileScan csv [order_id#0,franchise_id#1,store_id#2,product_id#3,quantity#4,unit_price#5,order_ts#6] Batched: false, DataFilters: [isnotnull(store_id#2)], Format: CSV, ...
+- Sort [store_id#7 ASC NULLS FIRST], false, 0
+- Exchange hashpartitioning(store_id#7, 200), ENSURE_REQUIREMENTS, [plan_id=332]
+- Filter isnotnull(store_id#7)
+- FileScan csv [store_id#7,store_name#8,city#9] Batched: false, DataFilters: [isnotnull(store_id#7)], Format: CSV, ...
joined_forced.count() = 10000000
Verificacion: SortMergeJoin forzado, mismo resultado, plan distinto -> OK
spark.sql.autoBroadcastJoinThreshold restaurado = 10485760
Parte 5 (L5-L6) -- revenue acumulado por tienda, con tiebreaker
running_total final por tienda = {'S02': 9700000.0, 'S01': 9575000.0, 'S03': 7262500.0}
Verificacion: running_total final coincide con el desglose conocido -> OK
Parte 6 (L7) -- producto top por tienda y por dia
num_top (filas con rank == 1) = 20
S01 2026-08-08 -- producto top = P004, total_revenue = 2250000.0
Verificacion: 20 combinaciones store x dia, cada una con su producto top, revenue escalado x250,000 -> OK
Parte 7 -- resumen final
Filas procesadas: 10,000,000
Revenue total verificado: 26,537,500.0
=== spark.stop() -- modulo 5 cerrado, joins y ventanas verificados sobre 10M filas ===
Siete partes, siete verificaciones, y el mismo 26,537,500.00 que el módulo 4 ya había confirmado con assert sobre los datos crudos —ahora reproducido de nuevo, esta vez sumando los tres running_total finales de la lección 6, en vez de un groupBy directo—. Fíjate en que la restauración del umbral, al final de la Parte 4 (10485760, sin la letra b esta vez, porque se estableció con un entero en vez de leerse del valor por defecto de Spark), confirma una disciplina importante para cualquier script real: cualquier configuración que cambies a propósito para un experimento debe volver a su valor original antes de que el resto del script dependa del comportamiento por defecto.
Diagrama: las siete partes, cerrando el módulo completo
flowchart TD
A["Parte 1: SparkSession abierta"] --> B
B["Parte 2 (L2-L4): fact_orders_at_scale,\ndim_store, dim_product leidas -- 10M filas"] --> C
C["Parte 3 (L2-L3): join por defecto --\nBroadcastHashJoin, sin Exchange de shuffle"] --> D
D["Parte 4 (L4): mismo join forzado --\nSortMergeJoin, Exchange en ambos lados"] --> E
E["Parte 5 (L5-L6): revenue acumulado\npor tienda, con desempate -- 9.57M/9.7M/7.26M"] --> F
F["Parte 6 (L7): producto top por\ntienda y dia -- 20 combinaciones"] --> G["Modulo 5 cerrado:\njoins y ventanas\nverificados sobre 10M filas"]
Cerrando el checklist de este módulo, pieza por pieza
| Pieza del módulo | Estado al cerrar este proyecto |
|---|---|
El criterio completo de broadcast join (autoBroadcastJoinThreshold) | Resuelto — lección 2, 10 MB por defecto, citado de la documentación oficial |
BroadcastHashJoin leído en .explain(), sobre fact_orders_at_scale | Resuelto — lección 3, cero Exchange de shuffle, verificado sobre 10M filas |
SortMergeJoin forzado y leído en .explain(), mismo join | Resuelto — lección 4, Exchange en ambos lados, mismo resultado |
Window.partitionBy().orderBy() — sintaxis básica | Resuelto — lección 5, verificado sobre las 40 filas reales |
| Revenue acumulado por tienda, con desempate a escala | Resuelto — lección 6, 9,575,000.00 / 9,700,000.00 / 7,262,500.00 |
| Ranking de producto top por tienda y por día | Resuelto — lección 7, 20 combinaciones, idénticas a 40 filas y a escala |
| Catalyst en sus fases, AQE, caché con criterio | Pendiente — módulo 6 |
| Parquet particionado a escala, UDFs vectorizados | Pendiente — módulo 7 |
| Capstone distribuido, árbol de decisión completo | Pendiente — módulo 8 |
Cinco módulos resueltos de ocho — más de la mitad del camino de esta guía. Con este módulo cerrado, tienes las dos piezas que faltaban para razonar con criterio completo sobre cualquier consulta de Spark: cuándo un JOIN paga el costo de un shuffle y cuándo no, y cómo responder preguntas de acumulado y ranking sin perder el detalle de fila que un groupBy sacrifica.
Errores comunes
Correr este proyecto sin haber generado kiosko_orders_at_scale.csv, dim_store.csv o dim_product.csv primero. Qué pasa: alguien salta directo a este proyecto sin tener los tres archivos en la carpeta de trabajo, y el script falla en la Parte 2 con un error de archivo no encontrado. Por qué pasa: es tentador tratar el proyecto de cierre como un punto de partida independiente. Cómo detectarlo: si spark.read.csv(...) falla con AnalysisException: Path does not exist, te falta alguno de los tres archivos —kiosko_orders_at_scale.csv del módulo 4, dim_store.csv/dim_product.csv del módulo 3— en la misma carpeta donde corres este script. Cómo corregirlo: este proyecto, igual que los de los módulos 3 y 4, reutiliza artefactos generados en lecciones anteriores — confirma que los tres archivos existen antes de correr el script completo.
Olvidar restaurar spark.sql.autoBroadcastJoinThreshold entre la Parte 4 y la Parte 5. Qué pasa: alguien elimina la línea que restaura el umbral al final de la Parte 4, y no nota ningún problema inmediato porque las Partes 5 y 6 de este proyecto no vuelven a usar .join() — pero si extendiera el script con un JOIN adicional después, ese JOIN heredaría el umbral desactivado sin ninguna advertencia. Por qué pasa: dentro de este proyecto específico, el efecto de omitir la restauración es invisible, porque ninguna parte posterior depende del comportamiento por defecto del broadcast. Cómo detectarlo: revisa que cualquier script que combine, en la misma SparkSession, un experimento de configuración (como forzar SortMergeJoin) con trabajo posterior que debería comportarse "normalmente", restaure explícitamente esa configuración. Cómo corregirlo: la Parte 4 de este proyecto restaura el umbral explícitamente, con un comentario claro de qué está haciendo y por qué — sigue ese mismo patrón en cualquier script propio que module configuraciones de Spark a mitad de camino.
Entregar el proyecto sin los assert de la Parte 6, confiando solo en .show(). Qué pasa: alguien corre las siete partes, ve que top_products.show() "se ve bien", y considera el proyecto terminado sin revisar si el assert de sample["product_id"] == "P004" realmente pasó. Por qué pasa: veinte filas con un rank de 1 en cada una parecen correctas a simple vista, y verificar formalmente un caso puntual parece un paso extra sobre un resultado que ya se ve bien. Cómo detectarlo: si tu entrega final no ejecutó los assert de la Parte 6 hasta el final sin lanzar AssertionError, no tienes ninguna garantía real de que el ranking sea correcto — la misma trampa que ya advirtieron los proyectos de los módulos 3 y 4. Cómo corregirlo: los assert de este proyecto —sobre el conteo de filas, el desglose por tienda, el producto top de una combinación específica— no son decorativos: son la prueba de que el script completo, no solo un fragmento, produce el resultado correcto de punta a punta.
Ejercicios
Ejercicio 1 — Extiende el proyecto con una Parte 8: el ranking, con dim_product unido para mostrar product_name en vez de product_id. Agrega una octava parte que una top_products contra dim_product_df (usando BroadcastHashJoin por defecto), y muestre el product_name legible en vez del código.
Ver solución
print("Parte 8 -- producto top, con nombre legible")
top_products_named = top_products.join(dim_product_df, "product_id").select(
"store_id", "order_day", "product_id", "product_name", "total_revenue", "rank"
).orderBy("store_id", "order_day")
top_products_named.explain()
top_products_named.show(20, truncate=False)
num_named = top_products_named.count()
assert num_named == 20
print("Verificacion: el join contra dim_product no perdio ninguna de las 20 filas -> OK")
Salida esperada (fragmento, ejecutada en esta corrida):
Parte 8 -- producto top, con nombre legible
== Physical Plan ==
AdaptiveSparkPlan isFinalPlan=false
+- ...
+- BroadcastHashJoin [product_id#N], [product_id#M], Inner, BuildRight, false, false
...
+--------+----------+----------+--------------------+-------------+----+
|store_id|order_day |product_id|product_name |total_revenue|rank|
+--------+----------+----------+--------------------+-------------+----+
|S01 |2026-08-03|P004 |Phone Charger Cable |1125000.0 |1 |
|S01 |2026-08-04|P003 |Instant Coffee Sachet|562500.0 |1 |
...
Verificacion: el join contra dim_product no perdio ninguna de las 20 filas -> OK
Confirmado: top_products (veinte filas, ya un resultado agregado y chico) sigue calificando de sobra para un BroadcastHashJoin contra dim_product, exactamente el mismo criterio que ya viste en las lecciones 2 a 4 — no importa que el JOIN ocurra al final de un pipeline de ventanas, el criterio de tamaño en bytes sigue siendo el mismo.
Ejercicio 2 — Confirma que el orden de las Partes 3/4 (joins) y las Partes 5/6 (ventanas) no afecta el resultado final. Sin ejecutar nada, explica en 2-3 frases si reordenar el script —calculando el revenue acumulado y el ranking (Partes 5 y 6) antes de los experimentos de JOIN (Partes 3 y 4)— cambiaría alguno de los assert finales.
Ver solución
No cambiaría ningún resultado: las Partes 3 y 4 operan sobre joined_default/joined_forced (fact_orders_at_scale_df unida contra las dimensiones), mientras que las Partes 5 y 6 operan directamente sobre fact_orders_at_scale_df y revenue_by_product_day (una agregación propia, sin necesitar ningún JOIN contra dimensiones) — los cuatro flujos son independientes entre sí en términos de qué datos consumen, así que reordenarlos no cambiaría ningún assert. Sin embargo, igual que en los proyectos de los módulos 3 y 4, el orden del script no es arbitrario desde el punto de vista pedagógico: seguir el mismo orden en que construiste el conocimiento a lo largo del módulo —primero el criterio de JOIN, después las ventanas— hace que la narrativa del proyecto tenga sentido para alguien que lo lee de arriba hacia abajo por primera vez.
Ejercicio 3 — Explica, sin mirar el diseño de la guía, qué le falta a este pipeline para ser el capstone completo del módulo 8. En un párrafo de 4-6 frases, describe qué piezas le faltan a kiosko_scaled_joins_and_rankings.py para convertirse en el pipeline distribuido completo del módulo 8.
Ver solución
Hoy, este script lee fact_orders_at_scale, decide con criterio entre BroadcastHashJoin y SortMergeJoin, calcula un acumulado por tienda y un ranking de producto top —pero nunca lee ni escribe ningún plan de Catalyst en sus fases completas (parsed, analyzed, optimized, physical), nunca decide explícitamente cuándo cachear un DataFrame reusado en varias consultas (aquí, fact_orders_at_scale_df se reusa varias veces sin ningún .cache() explícito), y nunca escribe nada a Parquet particionado por store_id, ni usa un pandas_udf vectorizado para ninguna columna derivada. Le falta, además, cualquier comparación de Adaptive Query Execution activo contra desactivado sobre este pipeline específico (módulo 6), y la disciplina completa de leer .explain(mode="formatted") en sus fases separadas, en vez del modo por defecto usado en todo este módulo. El capstone del módulo 8 ensambla todas esas piezas —joins con criterio, ventanas, caché con criterio, Parquet particionado, un pandas_udf— en un solo pipeline que corre de punta a punta sobre los diez millones de filas, y cierra con el árbol de decisión completo de "¿necesito Spark?", aplicado tanto al Kiosko real de cuarenta filas como a un Kiosko hipotético mucho más grande.
Resumen y siguiente paso: el final del módulo 5
Con este mini-proyecto cierras el módulo 5 completo. Confirmaste, con .explain() ejecutado sobre fact_orders_at_scale, que Spark elige BroadcastHashJoin por defecto contra dim_store y dim_product —sin ningún Exchange de shuffle— y que forzar spark.sql.autoBroadcastJoinThreshold = -1 produce SortMergeJoin, con Exchange en ambos lados del JOIN, sobre el mismo resultado lógico. Calculaste el revenue acumulado por tienda con Window.partitionBy("store_id").orderBy("order_ts", "franchise_id", "order_id"), verificado a mano sobre las cuarenta filas reales y confirmado a escala completa (9,575,000.00 / 9,700,000.00 / 7,262,500.00). Y cerraste con el ranking de producto top por tienda y por día, con F.row_number() sobre datos ya agregados, confirmando que el mismo producto gana en cada combinación, sin importar si el dataset tiene cuarenta filas o diez millones.
Diste el paso central de este módulo: dejaste de tratar "estrategia de JOIN" como una caja negra que Spark decide por su cuenta, y aprendiste el criterio exacto que la gobierna; y dejaste de pensar en groupBy como la única herramienta para responder preguntas agregadas, incorporando las funciones de ventana como la herramienta correcta cuando el detalle de fila importa tanto como el agregado.
Hacia dónde sigues. El módulo 6 —catalyst-explain-and-caching— toma el mismo tipo de plan de .explain() que ya leíste en este módulo, y lo desarrolla en profundidad: las fases completas del optimizador Catalyst (parsed, analyzed, optimized, physical), Adaptive Query Execution comparado activo contra desactivado, y el criterio completo de cuándo .cache() ayuda de verdad y cuándo solo gasta memoria sin ningún beneficio.
Recursos
- Apache Spark — SQL Performance Tuning (Catalyst, estrategias de
JOIN,spark.sql.autoBroadcastJoinThreshold— la referencia central de las Partes 3 y 4 de este proyecto). spark.apache.org/docs/latest/sql-performance-tuning.html. En inglés. - PySpark —
pyspark.sql.Window(la referencia departitionBy/orderBy, base de las Partes 5 y 6 de este proyecto). spark.apache.org/docs/latest/api/python/reference/pyspark.sql/window.html. En inglés. - DISEÑO de esta guía (
spark-and-distributed-processing-guide/DISENO.md) — el objetivo completo del módulo 5 y su lugar en el plan de ocho módulos. - DISEÑO de
data-modeling-for-analytics-guide— la fuente de la pregunta de acumulado y ranking que este módulo resolvió con la ventana nativa de Spark.src/guides/data-modeling-for-analytics-guide/DISENO.md