ML Engineer MasterClass (October) | 4 seats left

Spark · Pivot tables
AmazonAmazon Analytics
Amazon · Aggregate and Analyze

Pivot tables

Compare products across line quantities using explicit pivot columns.

Step 1 of 6 · Learn

Choose pivot columns

groupBy defines one output row per product. pivot turns quantity values into columns; the explicit list fixes column order and avoids a separate discovery of values. This example counts lines containing one or two units. These are line counts, not units sold. Zero fills absent combinations for this report.

Lesson reference: PySpark and Scala

Choose pivot columns

groupBy defines one output row per product. pivot turns quantity values into columns; the explicit list fixes column order and avoids a separate discovery of values. This example counts lines containing one or two units. These are line counts, not units sold. Zero fills absent combinations for this report.

PySpark example

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

result = order_items.groupBy("product_id").pivot(
    "quantity", [1, 2]
).agg(count("*")).na.fill(0).orderBy("product_id")
result.show()

Scala example

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

val result = order_items.groupBy("product_id").pivot(
  "quantity", Seq(1, 2)
).agg(count("*")).na.fill(0).orderBy("product_id")
result.show()

Change the metric

The layout and metric are separate choices. Keep the quantity columns, but switch from counting lines to summing line_total. The total already covers all units on a line; multiplying by quantity again would overstate the value. Zero here represents an absent combination, not a general rule for missing measurements.

PySpark example

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

result = order_items.groupBy("product_id").pivot(
    "quantity", [1, 2, 3]
).agg(count("*")).na.fill(0).orderBy("product_id")
result.show()

Scala example

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

val result = order_items.groupBy("product_id").pivot(
  "quantity", Seq(1, 2, 3)
).agg(count("*")).na.fill(0).orderBy("product_id")
result.show()

Filter before pivoting

A source filter changes contributions and can remove whole product rows. The example includes every line. Filtering line_total before groupBy asks a different question from filtering one of the finished quantity columns.

PySpark example

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

result = order_items.groupBy("product_id").pivot(
    "quantity", [1, 2, 3]
).agg(sum("line_total")).na.fill(0).orderBy("product_id")
result.show()

Scala example

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

val result = order_items.groupBy("product_id").pivot(
  "quantity", Seq(1, 2, 3)
).agg(sum("line_total")).na.fill(0).orderBy("product_id")
result.show()
example.pyPySpark
1. Read order linesNine lines reference six catalog products.
2. Choose pivot columnsGroup by product and pivot the explicit quantity values.
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
Choose pivot columns
Result6 rows
product_id12
20111
20210
20311
20410
20510
20610
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.