Lesson reference: PySpark and Scala
Count every match
An inner join emits one row per matching pair. The example lookup has one row each for 101 and 102. Adding a second 101 row makes each order for that customer appear twice. Two left rows times two right matches produces four rows for customer 101.
PySpark example
from pyspark.sql.functions import col, count
customers = spark.createDataFrame(
[(101,"Ari"),(102,"Bo")],
"customer_id long, customer_name string"
)
result = orders.join(customers, "customer_id", "inner").select("order_id", "customer_name").orderBy("order_id", "customer_name")
result.show()
Scala example
import org.apache.spark.sql.functions._
import spark.implicits._
val customers = Seq(
(101L, "Ari"),
(102L, "Bo")
).toDF("customer_id", "customer_name")
val lookup = customers
val result = orders.join(lookup, Seq("customer_id"), "inner").select("order_id", "customer_name").orderBy("order_id", "customer_name")
result.show()
Inspect match counts
Group the joined rows by order_id to see how many matches each order produced. The lookup has two different customer records for 101, so orders 1001 and 1003 each have two matches. This diagnostic identifies multiplication rather than silently removing it.
PySpark example
from pyspark.sql.functions import col, count
customers = spark.createDataFrame(
[(101,"Ari"),(101,"Ari updated"),(102,"Bo")],
"customer_id long, customer_name string"
)
result = orders.join(customers, "customer_id", "inner").groupBy("order_id").agg(
count("*").alias("match_count")
).filter(col("match_count") >= 1).orderBy("order_id")
result.show()
Scala example
import org.apache.spark.sql.functions._
import spark.implicits._
val customers = Seq(
(101L, "Ari"),
(101L, "Ari updated"),
(102L, "Bo")
).toDF("customer_id", "customer_name")
val result = orders.join(customers, Seq("customer_id"), "inner").groupBy("order_id").agg(
count("*").alias("match_count")
).filter(col("match_count") >= 1).orderBy("order_id")
result.show()
Remove exact duplicates
This lookup repeats the identical (101, Ari) row. distinct on the lookup removes that exact duplicate before joining. Do not use key-only deduplication to choose between conflicting names such as Ari and Ari updated: the surviving record would not be a defined business decision.
PySpark example
from pyspark.sql.functions import col, count
customers = spark.createDataFrame(
[(101,"Ari"),(101,"Ari"),(102,"Bo")],
"customer_id long, customer_name string"
)
result = orders.join(customers, "customer_id", "inner").select("order_id", "customer_name").orderBy("order_id", "customer_name")
result.show()
Scala example
import org.apache.spark.sql.functions._
import spark.implicits._
val customers = Seq(
(101L, "Ari"),
(101L, "Ari"),
(102L, "Bo")
).toDF("customer_id", "customer_name")
val lookup = customers
val result = orders.join(lookup, Seq("customer_id"), "inner").select("order_id", "customer_name").orderBy("order_id", "customer_name")
result.show()