ML Engineer MasterClass (October) | 4 seats left

Spark · Multiple grouping levels
AmazonAmazon Analytics
Amazon · Aggregate and Analyze

Multiple grouping levels

Build product and order subtotals with rollup and cube.

Step 1 of 6 · Learn

Add a grand total

rollup with one key returns detail groups plus a grand total. Spark uses null for rolled-up keys; these source IDs are non-null, so ALL is safe as a display label. For nullable keys use grouping or grouping_id to distinguish actual nulls from subtotals. Display IDs are cast to strings to accommodate ALL.

Lesson reference: PySpark and Scala

Add a grand total

rollup with one key returns detail groups plus a grand total. Spark uses null for rolled-up keys; these source IDs are non-null, so ALL is safe as a display label. For nullable keys use grouping or grouping_id to distinguish actual nulls from subtotals. Display IDs are cast to strings to accommodate ALL.

PySpark example

from pyspark.sql.functions import col, count, sum, coalesce, lit, concat

result = order_items.rollup("product_id").agg(
    count("*").alias("metric")
).select(
    coalesce(col("product_id").cast("string"), lit("ALL")).alias("product_id"),
    col("metric")
).orderBy("product_id")
result.show()

Scala example

import org.apache.spark.sql.functions._

val result = order_items.rollup("product_id").agg(
    count("*").alias("metric")
).select(
    coalesce(col("product_id").cast("string"), lit("ALL")).alias("product_id"),
    col("metric")
).orderBy("product_id")
result.show()

Choose a hierarchy

rollup(product_id, order_id) emits product/order details, product subtotals, and a grand total. It does not emit order-only subtotals. Reversing keys changes the hierarchy. Both display keys are strings, so ordering is lexical rather than numeric.

PySpark example

from pyspark.sql.functions import col, count, sum, coalesce, lit, concat

result = order_items.rollup("product_id", "order_id").agg(
    count("*").alias("metric")
).select(
    coalesce(col("product_id").cast("string"), lit("ALL")).alias("product_id"),
    coalesce(col("order_id").cast("string"), lit("ALL")).alias("order_id"),
    col("metric")
).orderBy("product_id", "order_id")
result.show()

Scala example

import org.apache.spark.sql.functions._

val result = order_items.rollup("product_id", "order_id").agg(
    count("*").alias("metric")
).select(
    coalesce(col("product_id").cast("string"), lit("ALL")).alias("product_id"),
    coalesce(col("order_id").cast("string"), lit("ALL")).alias("order_id"),
    col("metric")
).orderBy("product_id", "order_id")
result.show()

Include every combination

cube produces product/order details, product-only and order-only subtotals, and the grand total. Two keys produce four grouping sets, not four output rows. Every line contributes at several levels; adding all displayed metrics would count the same value repeatedly.

PySpark example

from pyspark.sql.functions import col, count, sum, coalesce, lit, concat

result = order_items.cube("product_id", "order_id").agg(
    count("*").alias("metric")
).select(
    coalesce(col("product_id").cast("string"), lit("ALL")).alias("product_id"),
    coalesce(col("order_id").cast("string"), lit("ALL")).alias("order_id"),
    col("metric")
).orderBy("product_id", "order_id")
result.show()

Scala example

import org.apache.spark.sql.functions._

val result = order_items.cube("product_id", "order_id").agg(
    count("*").alias("metric")
).select(
    coalesce(col("product_id").cast("string"), lit("ALL")).alias("product_id"),
    coalesce(col("order_id").cast("string"), lit("ALL")).alias("order_id"),
    col("metric")
).orderBy("product_id", "order_id")
result.show()
example.pyPySpark
1. Read order linesNine lines reference six catalog products.
2. Add a grand totalAggregate the stated grouping sets; keep detail and subtotal levels distinct.
3. Inspect the outputCheck the output keys and metric before comparing totals.
Source order_items9 rows
order_item_idorder_idproduct_idquantityline_total
11001201159.5
21001202130
31002203299
41002206150
51003204135
610042033149.97
71004205170.02
81005203149.99
910062012120
Add a grand total
Result7 rows
product_idmetric
2012
2021
2033
2041
2051
2061
ALL9
Compare the source columns with the transformed result.

The example is loaded in the editor. Run it as written, then try a small change.

Runs on the Spark backend. First startup may take a moment.

Run your code to see the result.