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()