Módulo 3: Rebuilding Fact Orders With The Dataframe Api
Agrupando por tienda y por producto
Descripción
Con revenue ya calculado en la lección 4, esta lección hace la pregunta de negocio real: ¿cuánto vendió cada tienda? ¿Cuánto representó cada producto? .groupBy().agg() de la DataFrame API responde ambas preguntas — y, siguiendo el criterio que ya adelantó la lección 4, esta es la lección donde por fin aplicas F.round(), una sola vez, sobre el resultado ya sumado.
Conexión con el módulo. Esta lección produce los dos desgloses —por tienda y por producto— que la lección 6 va a comparar, número por número, contra los que ya calculaste con dict, DuckDB, Polars y SQL en las tres guías anteriores. Sin este .groupBy(), no hay nada que verificar.
Una analogía: separar las fichas ya armadas en montones, y sumar cada uno
Tienes, después de la lección 4, cuarenta fichas completas —cada una con tienda, producto, cantidad, precio y el subtotal ya calculado—. Agruparlas es, literalmente, separarlas en montones sobre una mesa: un montón por cada tienda, o un montón por cada producto, y después sumar el subtotal de cada montón por separado. .groupBy("store_id") hace el primer montón; .groupBy("product_id") hace el segundo; y .agg(F.sum("revenue")) es la suma de cada montón, hecha de una vez para los tres o los cuatro montones a la vez, en vez de sumarlos uno por uno a mano.
Ejemplo trabajado: dos agrupaciones, con y sin redondeo
Paso 1 — Retoma fact_orders_df de la lección 4
# group_by_store_and_product.py
from pyspark.sql import SparkSession
from pyspark.sql import functions as F
from pyspark.sql.functions import col
from pyspark.sql.types import (
StructType, StructField, StringType, IntegerType, DoubleType, TimestampType,
)
spark = SparkSession.builder.appName("kiosko-spark").master("local[*]").getOrCreate()
orders_schema = StructType([
StructField("order_id", StringType(), 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_df = spark.read.csv("orders_2026-08-*.csv", schema=orders_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_df = (
orders_df
.join(dim_store_df, "store_id")
.join(dim_product_df, "product_id")
.withColumn("revenue", col("quantity") * col("unit_price"))
)
Paso 2 — Primero, sin redondear (el problema que la lección 4 anticipó)
print("=== groupBy SIN F.round() ===")
fact_orders_df.groupBy("store_id").agg(F.sum("revenue").alias("total_revenue")).orderBy("store_id").show()
Qué esperar (ejecutado en esta corrida):
=== groupBy SIN F.round() ===
+--------+------------------+
|store_id| total_revenue|
+--------+------------------+
| S01| 38.3|
| S02| 38.8|
| S03|29.049999999999997|
+--------+------------------+
Ahí está, exactamente como anticipó la lección 4: S03 muestra 29.049999999999997, no 29.05 — el mismo artefacto de punto flotante, ahora visible después de sumar dieciséis valores, no en una sola fila. S01 y S02 dan valores limpios en este caso particular (por cómo se acomodan sus sumas binarias), pero no cuentes con esa suerte en general.
Paso 3 — Con F.round(), aplicado una sola vez sobre la suma
print("=== groupBy CON F.round() ===")
revenue_by_store = (
fact_orders_df
.groupBy("store_id")
.agg(F.round(F.sum("revenue"), 2).alias("total_revenue"))
.orderBy("store_id")
)
revenue_by_store.show()
print("=== Por producto, con conteo de ordenes y unidades ===")
revenue_by_product = (
fact_orders_df
.groupBy("product_id")
.agg(
F.round(F.sum("revenue"), 2).alias("total_revenue"),
F.count("*").alias("num_orders"),
F.sum("quantity").alias("total_units"),
)
.orderBy("product_id")
)
revenue_by_product.show()
spark.stop()
Qué esperar (ejecutado en esta corrida):
=== groupBy CON F.round() ===
+--------+-------------+
|store_id|total_revenue|
+--------+-------------+
| S01| 38.3|
| S02| 38.8|
| S03| 29.05|
+--------+-------------+
=== Por producto, con conteo de ordenes y unidades ===
+----------+-------------+----------+-----------+
|product_id|total_revenue|num_orders|total_units|
+----------+-------------+----------+-----------+
| P001| 33.55| 16| 61|
| P002| 21.6| 10| 18|
| P003| 10.5| 7| 14|
| P004| 40.5| 7| 9|
+----------+-------------+----------+-----------+
F.round(F.sum("revenue"), 2) corrige exactamente el mismo caso: S03 ahora muestra 29.05, limpio. Y el desglose por producto trae, además del revenue, dos columnas que ya conoces de otro ángulo: num_orders —el conteo de líneas de orden por producto (P001: 16, P002: 10, P003: 7, P004: 7)— es exactamente el mismo conteo que ya viste en el proyecto del módulo 1, ahora derivado de un groupBy distinto en vez de un conteo aparte. La consistencia entre esos dos números, calculados en dos lecciones distintas con dos consultas distintas, es evidencia adicional de que ambos son correctos.
Diagrama: dos agrupaciones distintas, sobre las mismas cuarenta filas
flowchart TD
A["fact_orders_df\n40 filas, revenue calculado"]
A -->|"groupBy('store_id')\n.agg(F.round(F.sum('revenue'), 2))"| B["Por tienda:\nS01=38.3, S02=38.8, S03=29.05"]
A -->|"groupBy('product_id')\n.agg(F.round(F.sum('revenue'), 2))"| C["Por producto:\nP001=33.55, P002=21.6,\nP003=10.5, P004=40.5"]
B -.->|"suma"| D["106.15"]
C -.->|"suma"| D
Profundización: F.sum() es una función de agregación, no una Column normal
Fíjate en la diferencia entre col("quantity") (usada en la lección 4) y F.sum("revenue") (usada aquí): la primera se refiere a una columna que ya existe, fila por fila; la segunda es una función de agregación — solo tiene sentido dentro de un .agg() que sigue a un .groupBy() (o sobre el DataFrame completo, sin agrupar, como viste en la comparación de la lección 4). pyspark.sql.functions (importado casi siempre como F, la convención que usa esta guía desde esta lección en adelante) trae docenas de estas funciones — F.sum(), F.avg(), F.count(), F.max(), F.min(), y muchas más— todas diseñadas para trabajar dentro de un .agg(). La sintaxis se parece, a propósito, a GROUP BY de SQL: .groupBy("store_id") es el GROUP BY store_id, y cada función dentro de .agg() es una columna del SELECT —F.sum("revenue") es, ni más ni menos, SUM(revenue)—.
Vale la pena notar también que puedes agrupar por más de una columna a la vez, si la pregunta de negocio lo requiere —por ejemplo, .groupBy("store_id", "product_id") respondería "¿cuánto vendió cada producto, en cada tienda?", una granularidad más fina que cualquiera de los dos groupBy de esta lección. Esta guía no necesita esa granularidad para verificar el total de 106.15, pero la sintaxis es idéntica: agregar una columna más a .groupBy().
Errores comunes
Redondear la columna revenue antes de agregar, en vez de redondear la suma. Qué pasa: alguien, queriendo evitar el problema de punto flotante desde antes, escribe F.sum(F.round(col("revenue"), 2)) en vez de F.round(F.sum(col("revenue")), 2) — redondeando cada fila individual antes de sumarla, exactamente el patrón que la lección 4 desaconsejó. Por qué pasa: ambas expresiones parecen intercambiables a primera vista, porque usan las mismas dos funciones. Cómo detectarlo: compara los dos resultados sobre los mismos datos — en Kiosko, a esta escala, probablemente coincidan, pero el criterio de esta guía (y de las tres anteriores del ecosistema) es aplicar el redondeo una sola vez, al final, para evitar que el error se acumule de forma impredecible a mayor escala. Cómo corregirlo: el orden importa — agrega primero (F.sum("revenue")), redondea después (F.round(..., 2)), como hace el ejemplo trabajado de esta lección.
Olvidar .orderBy() después de .groupBy(), y confiar en que el resultado sale ordenado. Qué pasa: alguien mira el resultado de .groupBy("store_id").agg(...) y asume que las filas van a salir en un orden fijo (por ejemplo, S01, S02, S03), sin agregar .orderBy("store_id"). Por qué pasa: en este caso particular, con solo tres tiendas, el resultado sí suele salir ordenado — pero eso es una coincidencia de implementación, no una garantía, exactamente la misma lección que ya aprendiste en el módulo 2 sobre .show() sin .orderBy(). Cómo detectarlo: si tu código depende de que las filas de un groupBy salgan en un orden específico sin un .orderBy() explícito, tienes el mismo riesgo silencioso que ya viste antes. Cómo corregirlo: agrega siempre .orderBy() cuando el orden del resultado importa para comparar contra un valor esperado —exactamente lo que hace la lección 6 de este módulo—.
Confundir F.count("*") con F.sum("quantity"). Qué pasa: alguien usa F.count("*") (que cuenta filas, es decir, líneas de orden) cuando en realidad quería F.sum("quantity") (que suma las unidades vendidas), o al revés. Por qué pasa: ambas responden preguntas parecidas —"¿cuánto de este producto?"— pero a niveles distintos: número de órdenes contra número de unidades. Cómo detectarlo: en el ejemplo trabajado de esta lección, P001 tiene num_orders = 16 pero total_units = 61 —dieciséis líneas de orden distintas, pero sesenta y una botellas de agua vendidas en total, porque varias órdenes piden más de una unidad—. Si tu análisis mezcla estos dos números sin darse cuenta, las conclusiones de negocio van a estar mal. Cómo corregirlo: nombra las columnas agregadas con claridad (num_orders contra total_units, como en el ejemplo trabajado) para que la diferencia sea imposible de pasar por alto al leer el resultado.
Ejercicios
Ejercicio 1 — Agrupa por tienda y por producto a la vez. Usando .groupBy("store_id", "product_id"), calcula el revenue redondeado para cada combinación de tienda y producto. Confirma que sumar las cuatro filas de S01 (una por producto) da 38.3, el mismo total que ya conoces para esa tienda.
Ver solución
revenue_by_store_product = (
fact_orders_df
.groupBy("store_id", "product_id")
.agg(F.round(F.sum("revenue"), 2).alias("total_revenue"))
.orderBy("store_id", "product_id")
)
revenue_by_store_product.show(20)
Salida esperada (ejecutada en esta corrida):
+--------+----------+-------------+
|store_id|product_id|total_revenue|
+--------+----------+-------------+
| S01| P001| 11.0|
| S01| P002| 4.8|
| S01| P003| 4.5|
| S01| P004| 18.0|
| S02| P001| 10.45|
| S02| P002| 9.6|
| S02| P003| 5.25|
| S02| P004| 13.5|
| S03| P001| 12.1|
| S03| P002| 7.2|
| S03| P003| 0.75|
| S03| P004| 9.0|
+--------+----------+-------------+
S01: 11.0 + 4.8 + 4.5 + 18.0 = 38.3 — exacto, confirmando que el desglose más fino (tienda y producto a la vez) es consistente con el desglose por tienda sola de la lección.
Ejercicio 2 — Calcula el revenue por categoría, sin agregar ninguna columna nueva a fact_orders_df. Usando .groupBy("category") —una columna que ya viene de dim_product gracias al JOIN de la lección 3—, calcula el revenue total por categoría, y confirma que beverages es la categoría con más revenue.
Ver solución
revenue_by_category = (
fact_orders_df
.groupBy("category")
.agg(F.round(F.sum("revenue"), 2).alias("total_revenue"))
.orderBy(F.desc("total_revenue"))
)
revenue_by_category.show()
Salida esperada:
+-----------+-------------+
| category|total_revenue|
+-----------+-------------+
| beverages| 44.05|
|electronics| 40.5|
| snacks| 21.6|
+-----------+-------------+
beverages (P001 + P003 = 33.55 + 10.5 = 44.05) es, en efecto, la categoría con más revenue — el mismo resultado, con los mismos tres valores, que ya calculaste con SQL en data-modeling-for-analytics-guide. Fíjate en que esta consulta nunca tocó dim_store ni necesitó ningún JOIN nuevo: category ya estaba disponible en fact_orders_df desde la lección 3, porque el JOIN contra dim_product la trajo consigo.
Ejercicio 3 — Explica, sin código, por qué num_orders y total_units pueden ser distintos para el mismo producto. En 2-3 frases, usando los números reales de P001 (num_orders=16, total_units=61), explica por qué estos dos números no tienen que coincidir.
Ver solución
num_orders cuenta filas de fact_orders_df —cada fila es una línea de orden distinta, independientemente de cuántas unidades pidió—, mientras que total_units suma la columna quantity de esas mismas filas. Si cada orden pidiera exactamente una unidad, los dos números coincidirían siempre; pero como quantity varía por orden (algunas piden una botella de agua, otras piden cinco o seis, como viste en el módulo 1), la suma de unidades (61) termina siendo mayor que el conteo de órdenes (16). La diferencia entre ambos números es, en sí misma, información de negocio: dieciséis clientes distintos compraron agua embotellada esa semana, pero en promedio compraron más de tres botellas y media por orden.
Resumen y siguiente paso
En esta lección agrupaste fact_orders_df por tienda y por producto con .groupBy().agg(), y viste, con evidencia ejecutada, tanto el problema de punto flotante que la suma hereda de la lección 4 (29.049999999999997) como su solución correcta (F.round(F.sum(...), 2), aplicado una sola vez sobre el total agregado, no sobre cada fila). Los dos desgloses resultantes —S01=38.3, S02=38.8, S03=29.05 por tienda, P001=33.55, P002=21.6, P003=10.5, P004=40.5 por producto— son los números que la lección 6 va a verificar contra las tres guías anteriores del ecosistema.
Antes de avanzar deberías poder: escribir de memoria un .groupBy().agg(F.round(F.sum(...), 2)); explicar la diferencia entre F.count("*") y F.sum("quantity"); y explicar por qué el orden de redondear-y-sumar importa, con un ejemplo concreto de esta lección.
La lección 6 es la lección central de todo el módulo: toma estos dos desgloses y los compara, con assert, contra los números exactos de dict, DuckDB, Polars y SQL.
Recursos
- Apache Spark —
pyspark.sql.functions(el módulo completo de funciones de agregación, incluyendoF.sum,F.count,F.round, usadas en esta lección). spark.apache.org/docs/latest/api/python/reference/pyspark.sql/functions.html. En inglés. - Apache Spark —
DataFrame.groupBy(la referencia exacta de la API de agrupación). spark.apache.org/docs/latest/api/python/reference/pyspark.sql/api/pyspark.sql.DataFrame.groupBy.html. En inglés. - DISEÑO de
data-modeling-for-analytics-guide— la fuente del desglose por categoría (beverages=44.05,electronics=40.5,snacks=21.6) que el ejercicio 2 de esta lección reproduce con Spark.src/guides/data-modeling-for-analytics-guide/DISENO.md