ML Engineer MasterClass (October) | 4 seats left

Spark · Duplicates
AmazonAmazon Analytics
Amazon · Clean and Shape Data

Duplicates

Separate replayed order lines from legitimate repeat purchases and audit their counts.

Step 1 of 6 · Learn

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()
example.pyPySpark
1. Prepare replayed_itemsAppend three replayed records without changing order_items.
2. Remove identical linesApply the displayed expression to the prepared input.
3. Inspect the resultCompare row identity and counts, not just repeated product IDs.
Source replayed_items12 rows
order_item_idorder_idproduct_idquantityline_total
11001201159.5
21001202130
31002203299
41002206150
51003204135
610042033149.97
71004205170.02
81005203149.99
910062012120
11001201159.5
21001202130
11001201159.5
Remove identical lines
Result9 rows
order_item_idorder_id
11001
21001
31002
41002
51003
61004
71004
81005
91006
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.