In [2]:
import findspark
findspark.init()

from pyspark.sql import SparkSession
from pyspark.sql.types import StructType, StructField, StringType, IntegerType, DoubleType

# create spark session
spark = (
    SparkSession.builder.appName("working with string & dates").getOrCreate()
)

In [25]:
# Emp Data & Schema

emp_data = [
    ["001","101","John Doe","30","Male","50000","2015-01-01"],
    ["002","101","Jane Smith","25","Female","45000","2016-02-15"],
    ["003","102","Bob Brown","35","Male","55000","2014-05-01"],
    ["004","102","Alice Lee","28","Female","48000","2017-09-30"],
    ["005","103","Jack Chan","40","Male","60000","2013-04-01"],
    ["006","103","Jill Wong","32","","52000","2018-07-01"],
    ["007","101","James Johnson","42","Male","70000","2012-03-15"],
    ["008","102","Kate Kim","29","Female","51000","2019-10-01"],
    ["009","103","Tom Tan","33","Male","58000","2016-06-01"],
    ["010","104","Lisa Lee","27","Female","47000","2018-08-01"],
    ["011","104","David Park","38","Male","65000","2015-11-01"],
    ["012","105","Susan Chen","31","Female","54000","2017-02-15"],
    ["013","106","Brian Kim","45","Male","75000","2011-07-01"],
    ["014","107","Emily Lee","26","Female","46000","2019-01-01"],
    ["015","106","Michael Lee","37","Male","63000","2014-09-30"],
    ["016","107","Kelly Zhang","30","Female","49000","2018-04-01"],
    ["017","105","George Wang","34","","57000","2016-03-15"],
    ["018","104","Nancy Liu","29","Female","50000","2017-06-01"],
    ["019","103","Steven Chen","36","Male","62000","2015-08-01"],
    ["020","102","Grace Kim","32","Female","53000","2018-11-01"]
]

emp_schema = ["emp_id", "dept_id", "emp_name", "age", "gender", "salary", "hire_date"]

In [26]:
emp_df = spark.createDataFrame(data=emp_data, schema=emp_schema)

In [27]:
emp_df.show()

+------+-------+-------------+---+------+------+----------+
|emp_id|dept_id|     emp_name|age|gender|salary| hire_date|
+------+-------+-------------+---+------+------+----------+
|   001|    101|     John Doe| 30|  Male| 50000|2015-01-01|
|   002|    101|   Jane Smith| 25|Female| 45000|2016-02-15|
|   003|    102|    Bob Brown| 35|  Male| 55000|2014-05-01|
|   004|    102|    Alice Lee| 28|Female| 48000|2017-09-30|
|   005|    103|    Jack Chan| 40|  Male| 60000|2013-04-01|
|   006|    103|    Jill Wong| 32|      | 52000|2018-07-01|
|   007|    101|James Johnson| 42|  Male| 70000|2012-03-15|
|   008|    102|     Kate Kim| 29|Female| 51000|2019-10-01|
|   009|    103|      Tom Tan| 33|  Male| 58000|2016-06-01|
|   010|    104|     Lisa Lee| 27|Female| 47000|2018-08-01|
|   011|    104|   David Park| 38|  Male| 65000|2015-11-01|
|   012|    105|   Susan Chen| 31|Female| 54000|2017-02-15|
|   013|    106|    Brian Kim| 45|  Male| 75000|2011-07-01|
|   014|    107|    Emily Lee| 26|Female

In [28]:
emp_df.printSchema()

root
 |-- emp_id: string (nullable = true)
 |-- dept_id: string (nullable = true)
 |-- emp_name: string (nullable = true)
 |-- age: string (nullable = true)
 |-- gender: string (nullable = true)
 |-- salary: string (nullable = true)
 |-- hire_date: string (nullable = true)



In [29]:
# Case When
# select employee_id, name, age, salary, gender,
# case when gender = 'Male' then 'M' when gender = 'Female' then 'F' else null end as new_gender, hire_date from emp
from pyspark.sql.functions import when, col, expr

# using case When
emp_gender_fixed = emp_df.withColumn("new_gender", when(col("gender") == 'Male', 'M')
                                 .when(col("gender") == 'Female', 'F')
                                 .otherwise(None)
                                 )

# using 'expr' for case expression
emp_gender_fixed_1 = emp_df.withColumn("new_gender", expr("case when gender = 'Male' then 'M' when gender = 'Female' then 'F' else null end"))

In [30]:
emp_gender_fixed.show()

+------+-------+-------------+---+------+------+----------+----------+
|emp_id|dept_id|     emp_name|age|gender|salary| hire_date|new_gender|
+------+-------+-------------+---+------+------+----------+----------+
|   001|    101|     John Doe| 30|  Male| 50000|2015-01-01|         M|
|   002|    101|   Jane Smith| 25|Female| 45000|2016-02-15|         F|
|   003|    102|    Bob Brown| 35|  Male| 55000|2014-05-01|         M|
|   004|    102|    Alice Lee| 28|Female| 48000|2017-09-30|         F|
|   005|    103|    Jack Chan| 40|  Male| 60000|2013-04-01|         M|
|   006|    103|    Jill Wong| 32|      | 52000|2018-07-01|      NULL|
|   007|    101|James Johnson| 42|  Male| 70000|2012-03-15|         M|
|   008|    102|     Kate Kim| 29|Female| 51000|2019-10-01|         F|
|   009|    103|      Tom Tan| 33|  Male| 58000|2016-06-01|         M|
|   010|    104|     Lisa Lee| 27|Female| 47000|2018-08-01|         F|
|   011|    104|   David Park| 38|  Male| 65000|2015-11-01|         M|
|   01

In [31]:
emp_gender_fixed_1.show()

+------+-------+-------------+---+------+------+----------+----------+
|emp_id|dept_id|     emp_name|age|gender|salary| hire_date|new_gender|
+------+-------+-------------+---+------+------+----------+----------+
|   001|    101|     John Doe| 30|  Male| 50000|2015-01-01|         M|
|   002|    101|   Jane Smith| 25|Female| 45000|2016-02-15|         F|
|   003|    102|    Bob Brown| 35|  Male| 55000|2014-05-01|         M|
|   004|    102|    Alice Lee| 28|Female| 48000|2017-09-30|         F|
|   005|    103|    Jack Chan| 40|  Male| 60000|2013-04-01|         M|
|   006|    103|    Jill Wong| 32|      | 52000|2018-07-01|      NULL|
|   007|    101|James Johnson| 42|  Male| 70000|2012-03-15|         M|
|   008|    102|     Kate Kim| 29|Female| 51000|2019-10-01|         F|
|   009|    103|      Tom Tan| 33|  Male| 58000|2016-06-01|         M|
|   010|    104|     Lisa Lee| 27|Female| 47000|2018-08-01|         F|
|   011|    104|   David Park| 38|  Male| 65000|2015-11-01|         M|
|   01

In [32]:
# Replace in Strings
# select employee_id, emp_name, replace(emp_name, 'J', 'Z') as new_name, age, salary, gender, new_gender, hire_date from emp_gender_fixed
from pyspark.sql.functions import regexp_replace

emp_name_fixed = emp_gender_fixed.withColumn("new_name", regexp_replace(col("emp_name"), "J", "Z"))

In [33]:
emp_name_fixed.show()

+------+-------+-------------+---+------+------+----------+----------+-------------+
|emp_id|dept_id|     emp_name|age|gender|salary| hire_date|new_gender|     new_name|
+------+-------+-------------+---+------+------+----------+----------+-------------+
|   001|    101|     John Doe| 30|  Male| 50000|2015-01-01|         M|     Zohn Doe|
|   002|    101|   Jane Smith| 25|Female| 45000|2016-02-15|         F|   Zane Smith|
|   003|    102|    Bob Brown| 35|  Male| 55000|2014-05-01|         M|    Bob Brown|
|   004|    102|    Alice Lee| 28|Female| 48000|2017-09-30|         F|    Alice Lee|
|   005|    103|    Jack Chan| 40|  Male| 60000|2013-04-01|         M|    Zack Chan|
|   006|    103|    Jill Wong| 32|      | 52000|2018-07-01|      NULL|    Zill Wong|
|   007|    101|James Johnson| 42|  Male| 70000|2012-03-15|         M|Zames Zohnson|
|   008|    102|     Kate Kim| 29|Female| 51000|2019-10-01|         F|     Kate Kim|
|   009|    103|      Tom Tan| 33|  Male| 58000|2016-06-01|      

In [34]:
# Convert Date
# select *,  to_date(hire_date, 'YYYY-MM-DD') as hire_date from emp_name_fixed
from pyspark.sql.functions import to_date

emp_date_fix = emp_name_fixed.withColumn("hire_date", to_date(col("hire_date"), 'yyyy-MM-dd'))

In [35]:
emp_date_fix.show()

+------+-------+-------------+---+------+------+----------+----------+-------------+
|emp_id|dept_id|     emp_name|age|gender|salary| hire_date|new_gender|     new_name|
+------+-------+-------------+---+------+------+----------+----------+-------------+
|   001|    101|     John Doe| 30|  Male| 50000|2015-01-01|         M|     Zohn Doe|
|   002|    101|   Jane Smith| 25|Female| 45000|2016-02-15|         F|   Zane Smith|
|   003|    102|    Bob Brown| 35|  Male| 55000|2014-05-01|         M|    Bob Brown|
|   004|    102|    Alice Lee| 28|Female| 48000|2017-09-30|         F|    Alice Lee|
|   005|    103|    Jack Chan| 40|  Male| 60000|2013-04-01|         M|    Zack Chan|
|   006|    103|    Jill Wong| 32|      | 52000|2018-07-01|      NULL|    Zill Wong|
|   007|    101|James Johnson| 42|  Male| 70000|2012-03-15|         M|Zames Zohnson|
|   008|    102|     Kate Kim| 29|Female| 51000|2019-10-01|         F|     Kate Kim|
|   009|    103|      Tom Tan| 33|  Male| 58000|2016-06-01|      

In [36]:
emp_date_fix.printSchema()

root
 |-- emp_id: string (nullable = true)
 |-- dept_id: string (nullable = true)
 |-- emp_name: string (nullable = true)
 |-- age: string (nullable = true)
 |-- gender: string (nullable = true)
 |-- salary: string (nullable = true)
 |-- hire_date: date (nullable = true)
 |-- new_gender: string (nullable = true)
 |-- new_name: string (nullable = true)



In [37]:
# Add Date Columns
# Add current_date, current_timestamp, extract year from hire_date
from pyspark.sql.functions import current_date, current_timestamp

emp_dated = emp_date_fix.withColumn("date_now", current_date()).withColumn("timestamp_now", current_timestamp())

In [38]:
emp_dated.show()

+------+-------+-------------+---+------+------+----------+----------+-------------+----------+--------------------+
|emp_id|dept_id|     emp_name|age|gender|salary| hire_date|new_gender|     new_name|  date_now|       timestamp_now|
+------+-------+-------------+---+------+------+----------+----------+-------------+----------+--------------------+
|   001|    101|     John Doe| 30|  Male| 50000|2015-01-01|         M|     Zohn Doe|2025-09-29|2025-09-29 20:24:...|
|   002|    101|   Jane Smith| 25|Female| 45000|2016-02-15|         F|   Zane Smith|2025-09-29|2025-09-29 20:24:...|
|   003|    102|    Bob Brown| 35|  Male| 55000|2014-05-01|         M|    Bob Brown|2025-09-29|2025-09-29 20:24:...|
|   004|    102|    Alice Lee| 28|Female| 48000|2017-09-30|         F|    Alice Lee|2025-09-29|2025-09-29 20:24:...|
|   005|    103|    Jack Chan| 40|  Male| 60000|2013-04-01|         M|    Zack Chan|2025-09-29|2025-09-29 20:24:...|
|   006|    103|    Jill Wong| 32|      | 52000|2018-07-01|     

In [39]:
emp_dated.show(truncate=False)

+------+-------+-------------+---+------+------+----------+----------+-------------+----------+--------------------------+
|emp_id|dept_id|emp_name     |age|gender|salary|hire_date |new_gender|new_name     |date_now  |timestamp_now             |
+------+-------+-------------+---+------+------+----------+----------+-------------+----------+--------------------------+
|001   |101    |John Doe     |30 |Male  |50000 |2015-01-01|M         |Zohn Doe     |2025-09-29|2025-09-29 20:24:50.504598|
|002   |101    |Jane Smith   |25 |Female|45000 |2016-02-15|F         |Zane Smith   |2025-09-29|2025-09-29 20:24:50.504598|
|003   |102    |Bob Brown    |35 |Male  |55000 |2014-05-01|M         |Bob Brown    |2025-09-29|2025-09-29 20:24:50.504598|
|004   |102    |Alice Lee    |28 |Female|48000 |2017-09-30|F         |Alice Lee    |2025-09-29|2025-09-29 20:24:50.504598|
|005   |103    |Jack Chan    |40 |Male  |60000 |2013-04-01|M         |Zack Chan    |2025-09-29|2025-09-29 20:24:50.504598|
|006   |103    |

In [40]:
# Drop Null gender records
emp_1 = emp_dated.na.drop()

In [41]:
emp_1.show()

+------+-------+-------------+---+------+------+----------+----------+-------------+----------+--------------------+
|emp_id|dept_id|     emp_name|age|gender|salary| hire_date|new_gender|     new_name|  date_now|       timestamp_now|
+------+-------+-------------+---+------+------+----------+----------+-------------+----------+--------------------+
|   001|    101|     John Doe| 30|  Male| 50000|2015-01-01|         M|     Zohn Doe|2025-09-29|2025-09-29 20:25:...|
|   002|    101|   Jane Smith| 25|Female| 45000|2016-02-15|         F|   Zane Smith|2025-09-29|2025-09-29 20:25:...|
|   003|    102|    Bob Brown| 35|  Male| 55000|2014-05-01|         M|    Bob Brown|2025-09-29|2025-09-29 20:25:...|
|   004|    102|    Alice Lee| 28|Female| 48000|2017-09-30|         F|    Alice Lee|2025-09-29|2025-09-29 20:25:...|
|   005|    103|    Jack Chan| 40|  Male| 60000|2013-04-01|         M|    Zack Chan|2025-09-29|2025-09-29 20:25:...|
|   007|    101|James Johnson| 42|  Male| 70000|2012-03-15|     

In [42]:
# Fix Null values
# select *, nvl('new_gender', 'O') as new_gender from emp_dated
from pyspark.sql.functions import coalesce, lit

emp_null_df = emp_dated.withColumn("new_gender", coalesce(col("new_gender"), lit("O")))

In [43]:
emp_null_df.show()

+------+-------+-------------+---+------+------+----------+----------+-------------+----------+--------------------+
|emp_id|dept_id|     emp_name|age|gender|salary| hire_date|new_gender|     new_name|  date_now|       timestamp_now|
+------+-------+-------------+---+------+------+----------+----------+-------------+----------+--------------------+
|   001|    101|     John Doe| 30|  Male| 50000|2015-01-01|         M|     Zohn Doe|2025-09-29|2025-09-29 20:28:...|
|   002|    101|   Jane Smith| 25|Female| 45000|2016-02-15|         F|   Zane Smith|2025-09-29|2025-09-29 20:28:...|
|   003|    102|    Bob Brown| 35|  Male| 55000|2014-05-01|         M|    Bob Brown|2025-09-29|2025-09-29 20:28:...|
|   004|    102|    Alice Lee| 28|Female| 48000|2017-09-30|         F|    Alice Lee|2025-09-29|2025-09-29 20:28:...|
|   005|    103|    Jack Chan| 40|  Male| 60000|2013-04-01|         M|    Zack Chan|2025-09-29|2025-09-29 20:28:...|
|   006|    103|    Jill Wong| 32|      | 52000|2018-07-01|     

In [44]:
# Drop old columns and Fix new column names
emp_final = emp_null_df.drop("emp_name", "gender").withColumnRenamed("new_name", "name").withColumnRenamed("new_gender", "gender")

In [45]:
emp_final.show()

+------+-------+---+------+----------+------+-------------+----------+--------------------+
|emp_id|dept_id|age|salary| hire_date|gender|         name|  date_now|       timestamp_now|
+------+-------+---+------+----------+------+-------------+----------+--------------------+
|   001|    101| 30| 50000|2015-01-01|     M|     Zohn Doe|2025-09-29|2025-09-29 20:30:...|
|   002|    101| 25| 45000|2016-02-15|     F|   Zane Smith|2025-09-29|2025-09-29 20:30:...|
|   003|    102| 35| 55000|2014-05-01|     M|    Bob Brown|2025-09-29|2025-09-29 20:30:...|
|   004|    102| 28| 48000|2017-09-30|     F|    Alice Lee|2025-09-29|2025-09-29 20:30:...|
|   005|    103| 40| 60000|2013-04-01|     M|    Zack Chan|2025-09-29|2025-09-29 20:30:...|
|   006|    103| 32| 52000|2018-07-01|     O|    Zill Wong|2025-09-29|2025-09-29 20:30:...|
|   007|    101| 42| 70000|2012-03-15|     M|Zames Zohnson|2025-09-29|2025-09-29 20:30:...|
|   008|    102| 29| 51000|2019-10-01|     F|     Kate Kim|2025-09-29|2025-09-29

In [46]:
# Write data as CSV
emp_final.write.format("csv").save("data/output/3/emp.csv")

In [47]:
# Bonus TIP
# Convert date into String and extract date information
from pyspark.sql.functions import date_format

emp_fixed = emp_final.withColumn("date_string", date_format(col("hire_date"), "dd/MM/yyyy"))
emp_fixed_1 = emp_final.withColumn("date_year", date_format(col("timestamp_now"), "z"))

In [48]:
emp_fixed.show()

+------+-------+---+------+----------+------+-------------+----------+--------------------+-----------+
|emp_id|dept_id|age|salary| hire_date|gender|         name|  date_now|       timestamp_now|date_string|
+------+-------+---+------+----------+------+-------------+----------+--------------------+-----------+
|   001|    101| 30| 50000|2015-01-01|     M|     Zohn Doe|2025-09-29|2025-09-29 20:36:...| 01/01/2015|
|   002|    101| 25| 45000|2016-02-15|     F|   Zane Smith|2025-09-29|2025-09-29 20:36:...| 15/02/2016|
|   003|    102| 35| 55000|2014-05-01|     M|    Bob Brown|2025-09-29|2025-09-29 20:36:...| 01/05/2014|
|   004|    102| 28| 48000|2017-09-30|     F|    Alice Lee|2025-09-29|2025-09-29 20:36:...| 30/09/2017|
|   005|    103| 40| 60000|2013-04-01|     M|    Zack Chan|2025-09-29|2025-09-29 20:36:...| 01/04/2013|
|   006|    103| 32| 52000|2018-07-01|     O|    Zill Wong|2025-09-29|2025-09-29 20:36:...| 01/07/2018|
|   007|    101| 42| 70000|2012-03-15|     M|Zames Zohnson|2025-

In [49]:
emp_fixed_1.show()

+------+-------+---+------+----------+------+-------------+----------+--------------------+---------+
|emp_id|dept_id|age|salary| hire_date|gender|         name|  date_now|       timestamp_now|date_year|
+------+-------+---+------+----------+------+-------------+----------+--------------------+---------+
|   001|    101| 30| 50000|2015-01-01|     M|     Zohn Doe|2025-09-29|2025-09-29 20:36:...|      IST|
|   002|    101| 25| 45000|2016-02-15|     F|   Zane Smith|2025-09-29|2025-09-29 20:36:...|      IST|
|   003|    102| 35| 55000|2014-05-01|     M|    Bob Brown|2025-09-29|2025-09-29 20:36:...|      IST|
|   004|    102| 28| 48000|2017-09-30|     F|    Alice Lee|2025-09-29|2025-09-29 20:36:...|      IST|
|   005|    103| 40| 60000|2013-04-01|     M|    Zack Chan|2025-09-29|2025-09-29 20:36:...|      IST|
|   006|    103| 32| 52000|2018-07-01|     O|    Zill Wong|2025-09-29|2025-09-29 20:36:...|      IST|
|   007|    101| 42| 70000|2012-03-15|     M|Zames Zohnson|2025-09-29|2025-09-29 2

In [None]:
# doc links:
# https://spark.apache.org/docs/latest/sql-ref-datetime-pattern.html
# https://spark.apache.org/docs/latest/api/python/reference/pyspark.sql/functions.html

In [50]:
spark.stop()