-sandbox

<div style="text-align: center; line-height: 0; padding-top: 9px;">
  <img src="https://databricks.com/wp-content/uploads/2018/03/db-academy-rgb-1200px.png" alt="Databricks Learning" style="width: 600px">
</div>

# De-Duping Data
##![Spark Logo Tiny](https://files.training.databricks.com/images/105/logo_spark_tiny.png) Instructions

In this exercise, we're doing ETL on a file we've received from some customer. That file contains data about people, including:

* first, middle and last names
* gender
* birth date
* Social Security number
* salary

But, as is unfortunately common in data we get from this customer, the file contains some duplicate records. Worse:

* In some of the records, the names are mixed case (e.g., "Carol"), while in others, they are uppercase (e.g., "CAROL"). 
* The Social Security numbers aren't consistent, either. Some of them are hyphenated (e.g., "992-83-4829"), while others are missing hyphens ("992834829").

The name fields are guaranteed to match, if you disregard character case, and the birth dates will also match. (The salaries will match, as well,
and the Social Security Numbers *would* match, if they were somehow put in the same format).

Your job is to remove the duplicate records. The specific requirements of your job are:

* Remove duplicates. It doesn't matter which record you keep; it only matters that you keep one of them.
* Preserve the data format of the columns. For example, if you write the first name column in all lower-case, you haven't met this requirement.
* Write the result as a Parquet file, as designated by *dest_file*.
* The final Parquet "file" must contain 8 part files (8 files ending in ".parquet").

<img src="https://files.training.databricks.com/images/icon_hint_24.png"/>&nbsp;**Hint:** The initial dataset contains 103,000 records.<br/>
The de-duplicated result haves 100,000 records.

##![Spark Logo Tiny](https://files.training.databricks.com/images/105/logo_spark_tiny.png) Getting Started

Run the following cell to configure our "classroom."

In [0]:
!pip install mlflow

In [0]:
%run "../Includes/Classroom-Setup"

##![Spark Logo Tiny](https://files.training.databricks.com/images/105/logo_spark_tiny.png) Hints

* Use the <a href="https://spark.apache.org/docs/latest/api/python/index.html" target="_blank">API docs</a>. Specifically, you might find 
  <a href="https://spark.apache.org/docs/latest/api/python/reference/api/pyspark.sql.DataFrame.html?highlight=dataframe#pyspark.sql.DataFrame" target="_blank">DataFrame</a> and
  <a href="https://spark.apache.org/docs/latest/api/python/reference/pyspark.sql.html#functions" target="_blank">functions</a> to be helpful.
* It's helpful to look at the file first, so you can check the format. **`dbutils.fs.head()`** (or just **`%fs head`**) is a big help here.

In [0]:
# TODO

source_file = f"{datasets_dir}/dataframes/people-with-dups.txt"
dest_file = userhome + "/people.parquet"


# In case it already exists
dbutils.fs.rm(dest_file, True)

In [0]:
# ANSWER

# dropDuplicates() will likely introduce a shuffle, so it helps to reduce the number of post-shuffle partitions.
spark.conf.set("spark.sql.shuffle.partitions", 8)

In [0]:
# ANSWER

df = (spark
      .read
      .option("header", "true")
      .option("inferSchema", "true")
      .option("sep", ":")
      .csv(source_file))

df.count()

In [0]:
# ANSWER
from pyspark.sql.functions import  col, lower, translate

deduped_df = (df
              .select(col("*"),
                      lower(col("firstName")).alias("lcFirstName"),
                      lower(col("lastName")).alias("lcLastName"),
                      lower(col("middleName")).alias("lcMiddleName"),
                      translate(col("ssn"), "-", "").alias("ssnNums")
                      # regexp_replace(col("ssn"), "-", "").alias("ssnNums")
                      # regexp_replace(col("ssn"), """^(\d{3})(\d{2})(\d{4})$""", "$1-$2-$3").alias("ssnNums")
                     )
             .dropDuplicates(["lcFirstName", "lcMiddleName", "lcLastName", "ssnNums", "gender", "birthDate", "salary"])
             .drop("lcFirstName", "lcMiddleName", "lcLastName", "ssnNums")
            )

In [0]:
# ANSWER

# Now we can save the results. We'll also re-read them and count them, just as a final check.
# Just for fun, we'll use the Snappy compression codec. It's not as compact as Gzip, but it's much faster.
deduped_df.write.mode("overwrite").parquet(dest_file)

deduped_df = spark.read.parquet(dest_file)
print(f"Total Records: {deduped_df.count()}")

In [0]:
# ANSWER

display(dbutils.fs.ls(dest_file))

path,name,size,modificationTime
dbfs:/user/manujkumar.joshi@celebaltech.com/dbacademy/people.parquet/_SUCCESS,_SUCCESS,0,1661401876000
dbfs:/user/manujkumar.joshi@celebaltech.com/dbacademy/people.parquet/_committed_1240862853632513444,_committed_1240862853632513444,420,1661401875000
dbfs:/user/manujkumar.joshi@celebaltech.com/dbacademy/people.parquet/_started_1240862853632513444,_started_1240862853632513444,0,1661401873000
dbfs:/user/manujkumar.joshi@celebaltech.com/dbacademy/people.parquet/part-00000-tid-1240862853632513444-7b16f52d-69ef-49cf-90c7-c81727a7add4-69-1-c000.snappy.parquet,part-00000-tid-1240862853632513444-7b16f52d-69ef-49cf-90c7-c81727a7add4-69-1-c000.snappy.parquet,780865,1661401875000
dbfs:/user/manujkumar.joshi@celebaltech.com/dbacademy/people.parquet/part-00001-tid-1240862853632513444-7b16f52d-69ef-49cf-90c7-c81727a7add4-70-1-c000.snappy.parquet,part-00001-tid-1240862853632513444-7b16f52d-69ef-49cf-90c7-c81727a7add4-70-1-c000.snappy.parquet,775438,1661401875000
dbfs:/user/manujkumar.joshi@celebaltech.com/dbacademy/people.parquet/part-00002-tid-1240862853632513444-7b16f52d-69ef-49cf-90c7-c81727a7add4-71-1-c000.snappy.parquet,part-00002-tid-1240862853632513444-7b16f52d-69ef-49cf-90c7-c81727a7add4-71-1-c000.snappy.parquet,778067,1661401875000
dbfs:/user/manujkumar.joshi@celebaltech.com/dbacademy/people.parquet/part-00003-tid-1240862853632513444-7b16f52d-69ef-49cf-90c7-c81727a7add4-72-1-c000.snappy.parquet,part-00003-tid-1240862853632513444-7b16f52d-69ef-49cf-90c7-c81727a7add4-72-1-c000.snappy.parquet,769956,1661401875000


##![Spark Logo Tiny](https://s3-us-west-2.amazonaws.com/curriculum-release/images/105/logo_spark_tiny.png) Validate Your Answer

At the bare minimum, we can verify that you wrote the parquet file out to **dest_file** and that you have the right number of records.

Running the following cell to confirm your result:

In [0]:
part_files = len(list(filter(lambda f: f.path.endswith(".parquet"), dbutils.fs.ls(dest_file))))

final_df = spark.read.parquet(dest_file)
final_count = final_df.count()

clearYourResults()
validateYourAnswer("01 Parquet File Exists", 1276280174, part_files)
validateYourAnswer("02 Expected 100000 Records", 972882115, final_count)
summarizeYourResults()

0,1,2
01 Parquet File Exists:,FAILED,4
02 Expected 100000 Records:,passed,100000


-sandbox
&copy; 2022 Databricks, Inc. All rights reserved.<br/>
Apache, Apache Spark, Spark and the Spark logo are trademarks of the <a href="https://www.apache.org/">Apache Software Foundation</a>.<br/>
<br/>
<a href="https://databricks.com/privacy-policy">Privacy Policy</a> | <a href="https://databricks.com/terms-of-use">Terms of Use</a> | <a href="https://help.databricks.com/">Support</a>