Remove identical lines
The setup appends copies of order lines 1 and 2, then another copy of line 1: nine source lines become twelve records. distinct compares every selected column. The example selects order_item_id and order_id first; fewer selected columns can change what counts as identical.
Lesson reference: PySpark and Scala
Remove identical lines
The setup appends copies of order lines 1 and 2, then another copy of line 1: nine source lines become twelve records. distinct compares every selected column. The example selects order_item_id and order_id first; fewer selected columns can change what counts as identical.
PySpark example
from pyspark.sql.functions import (
col, lit, concat, trim, lower, upper, regexp_extract,
format_string, to_date, to_timestamp, datediff
)
replayed_items = (
order_items.unionByName(order_items.filter(col("order_item_id") <= 2))
.unionByName(order_items.filter(col("order_item_id") == 1))
)
result = replayed_items.select("order_item_id", "order_id").distinct().orderBy("order_item_id")
result.show()
Scala example
import org.apache.spark.sql.functions._
val replayed_items = order_items
.unionByName(order_items.filter(col("order_item_id") <= 2))
.unionByName(order_items.filter(col("order_item_id") === 1))
val result = replayed_items.select("order_item_id", "order_id").distinct().orderBy("order_item_id")
result.show()
Choose the identity
The same product can appear in different orders without being a duplicate line. The example enumerates product IDs. Adding order_id enumerates product/order combinations; it does not identify individual lines in a general dataset. Dropping duplicates by a key while retaining conflicting non-key values does not guarantee which record survives.
PySpark example
from pyspark.sql.functions import (
col, lit, concat, trim, lower, upper, regexp_extract,
format_string, to_date, to_timestamp, datediff
)
replayed_items = (
order_items.unionByName(order_items.filter(col("order_item_id") <= 2))
.unionByName(order_items.filter(col("order_item_id") == 1))
)
result = replayed_items.select("product_id").distinct().orderBy("product_id")
result.show()
Scala example
import org.apache.spark.sql.functions._
val replayed_items = order_items
.unionByName(order_items.filter(col("order_item_id") <= 2))
.unionByName(order_items.filter(col("order_item_id") === 1))
val result = replayed_items.select("product_id").distinct().orderBy("product_id")
result.show()
Audit replay counts
groupBy(order_item_id).count() counts occurrences of each line ID. The example keeps IDs occurring at least three times. Counting before deduplication exposes replays; it does not choose an authoritative record or sum purchased quantities.
PySpark example
from pyspark.sql.functions import (
col, lit, concat, trim, lower, upper, regexp_extract,
format_string, to_date, to_timestamp, datediff
)
replayed_items = (
order_items.unionByName(order_items.filter(col("order_item_id") <= 2))
.unionByName(order_items.filter(col("order_item_id") == 1))
)
result = replayed_items.groupBy("order_item_id").count().filter(col("count") >= 3).orderBy("order_item_id")
result.show()
Scala example
import org.apache.spark.sql.functions._
val replayed_items = order_items
.unionByName(order_items.filter(col("order_item_id") <= 2))
.unionByName(order_items.filter(col("order_item_id") === 1))
val result = replayed_items.groupBy("order_item_id").count().filter(col("count") >= 3).orderBy("order_item_id")
result.show()