ML Engineer MasterClass (October) | 4 seats left

Spark · Inner joins
AmazonAmazon Analytics
Amazon · Combine DataFrames

Inner joins

Enrich order lines with names and categories from the product catalog.

Step 1 of 6 · Learn

Match catalog keys

selected_products is a filtered view of the existing products table, not a replacement catalog. The example includes IDs up to 203. An inner join retains only lines with a match. The catalog has one row per product_id, so each matched line appears once; a key-name join keeps one product_id column.

Lesson reference: PySpark and Scala

Match catalog keys

selected_products is a filtered view of the existing products table, not a replacement catalog. The example includes IDs up to 203. An inner join retains only lines with a match. The catalog has one row per product_id, so each matched line appears once; a key-name join keeps one product_id column.

PySpark example

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

selected_products = products.filter(col("product_id") <= 203)

result = order_items.join(selected_products, "product_id", "inner").select("order_item_id", "product_id", "product_name").orderBy("order_item_id")
result.show()

Scala example

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

val selected_products = products.filter(col("product_id") <= 203)

val result = order_items.join(selected_products, Seq("product_id"), "inner").select("order_item_id", "product_id", "product_name").orderBy("order_item_id")
result.show()

Choose joined fields

The join combines line values from order_items with product names from the catalog. Select only the fields the report needs. Repeated purchases of a product are legitimate separate lines, not duplicate catalog records. If non-key names overlap, qualify them with DataFrame aliases.

PySpark example

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

selected_products = products.filter(col("product_id") <= 203)

result = order_items.join(selected_products, "product_id", "inner").select("order_item_id", "product_name", "line_total").orderBy("order_item_id")
result.show()

Scala example

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

val selected_products = products.filter(col("product_id") <= 203)

val result = order_items.join(selected_products, Seq("product_id"), "inner").select("order_item_id", "product_name", "line_total").orderBy("order_item_id")
result.show()

Filter matched lines

A category filter after the join uses a field from the catalog. The example keeps Electronics matches within the selected subset. A filter cannot restore lines whose products were excluded before the join; those matches never existed.

PySpark example

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

selected_products = products.filter(col("product_id") <= 203)

result = order_items.join(selected_products, "product_id", "inner").filter(col("category") == "Electronics").select("order_item_id", "product_name", "line_total", "category").orderBy("order_item_id")
result.show()

Scala example

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

val selected_products = products.filter(col("product_id") <= 203)

val result = order_items.join(selected_products, Seq("product_id"), "inner").filter(col("category") === "Electronics").select("order_item_id", "product_name", "line_total", "category").orderBy("order_item_id")
result.show()
example.pyPySpark
1. Read order linesNine lines reference six catalog products.
2. Match catalog keysMatch each line to one selected catalog row.
3. Inspect the outputUnmatched lines are absent; repeated purchases remain separate lines.
Source order_items9 rows
order_item_idorder_idproduct_idquantityline_total
11001201159.5
21001202130
31002203299
41002206150
51003204135
610042033149.97
71004205170.02
81005203149.99
910062012120
Second input selected_products3 rows
product_idproduct_namecategory
201Wireless keyboardElectronics
202Laptop standOffice
203Desk lampHome
Match catalog keys
Result6 rows
order_item_idproduct_idproduct_name
1201Wireless keyboard
2202Laptop stand
3203Desk lamp
6203Desk lamp
8203Desk lamp
9201Wireless keyboard
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.