In [1]:
from pyspark.sql import SparkSession
import pyspark.sql.functions as F
import os

In [2]:
os.environ['SPARK_HOME'] = r'C:\spark\spark-3.5.4-bin-hadoop3'
os.environ['PYSPARK_DRIVER_PYTHON'] = 'jupyter'
os.environ['PYSPARK_DRIVER_PYTHON_OPTS'] = 'lab'
os.environ['PYSPARK_PYTHON'] = 'python'

In [3]:
spark = (
    SparkSession
    .builder
    .appName("PySpark Zero to Hero")
    .master("local[*]")
    .config("spark.executor.memory", "16g")
    .config("spark.driver.memory", "16g")
    .config("spark.executor.cores", "4")
    .config("spark.sql.shuffle.partitions", "80")
    .config("spark.dynamicAllocation.enabled", "true")
    .config("spark.dynamicAllocation.minExecutors", "2")
    .config("spark.dynamicAllocation.initialExecutors", "24")
    .config("spark.dynamicAllocation.maxExecutors", "50")
    .config("spark.shuffle.service.enabled", "true")
    .config("spark.serializer", "org.apache.spark.serializer.KryoSerializer")
    .getOrCreate()
)

In [6]:
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","Female","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","Male","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"]
]

In [7]:
emp_schema = "employee_id string, department_id string, name string, age string, gender string, salary string, hire_date string"

In [8]:
emp = spark.createDataFrame(data=emp_data, schema=emp_schema)

In [9]:
emp_casted = emp.select(
    'employee_id', 'name', 'age', F.col('salary')
    .cast('double').alias('salary')
)

In [10]:
emp_casted.printSchema()

root
 |-- employee_id: string (nullable = true)
 |-- name: string (nullable = true)
 |-- age: string (nullable = true)
 |-- salary: double (nullable = true)



In [11]:
emp_taxed = emp_casted.withColumn(
    'tax', F.col('salary') * 0.2
)

In [12]:
emp_taxed.show()

+-----------+-------------+---+-------+-------+
|employee_id|         name|age| salary|    tax|
+-----------+-------------+---+-------+-------+
|        001|     John Doe| 30|50000.0|10000.0|
|        002|   Jane Smith| 25|45000.0| 9000.0|
|        003|    Bob Brown| 35|55000.0|11000.0|
|        004|    Alice Lee| 28|48000.0| 9600.0|
|        005|    Jack Chan| 40|60000.0|12000.0|
|        006|    Jill Wong| 32|52000.0|10400.0|
|        007|James Johnson| 42|70000.0|14000.0|
|        008|     Kate Kim| 29|51000.0|10200.0|
|        009|      Tom Tan| 33|58000.0|11600.0|
|        010|     Lisa Lee| 27|47000.0| 9400.0|
|        011|   David Park| 38|65000.0|13000.0|
|        012|   Susan Chen| 31|54000.0|10800.0|
|        013|    Brian Kim| 45|75000.0|15000.0|
|        014|    Emily Lee| 26|46000.0| 9200.0|
|        015|  Michael Lee| 37|63000.0|12600.0|
|        016|  Kelly Zhang| 30|49000.0| 9800.0|
|        017|  George Wang| 34|57000.0|11400.0|
|        018|    Nancy Liu| 29|50000.0|1

In [13]:
emp_new_cols = emp_taxed.withColumn(
    'ColumnOne', F.lit(1)
).withColumn(
    'ColumnTwo', F.lit('two')
)

In [14]:
emp_new_cols.show(5)

+-----------+----------+---+-------+-------+---------+---------+
|employee_id|      name|age| salary|    tax|ColumnOne|ColumnTwo|
+-----------+----------+---+-------+-------+---------+---------+
|        001|  John Doe| 30|50000.0|10000.0|        1|      two|
|        002|Jane Smith| 25|45000.0| 9000.0|        1|      two|
|        003| Bob Brown| 35|55000.0|11000.0|        1|      two|
|        004| Alice Lee| 28|48000.0| 9600.0|        1|      two|
|        005| Jack Chan| 40|60000.0|12000.0|        1|      two|
+-----------+----------+---+-------+-------+---------+---------+
only showing top 5 rows



In [15]:
emp_1 = emp_new_cols.withColumnRenamed(
    'employee_id', 'emp_id'
)

In [16]:
emp_1.show()

+------+-------------+---+-------+-------+---------+---------+
|emp_id|         name|age| salary|    tax|ColumnOne|ColumnTwo|
+------+-------------+---+-------+-------+---------+---------+
|   001|     John Doe| 30|50000.0|10000.0|        1|      two|
|   002|   Jane Smith| 25|45000.0| 9000.0|        1|      two|
|   003|    Bob Brown| 35|55000.0|11000.0|        1|      two|
|   004|    Alice Lee| 28|48000.0| 9600.0|        1|      two|
|   005|    Jack Chan| 40|60000.0|12000.0|        1|      two|
|   006|    Jill Wong| 32|52000.0|10400.0|        1|      two|
|   007|James Johnson| 42|70000.0|14000.0|        1|      two|
|   008|     Kate Kim| 29|51000.0|10200.0|        1|      two|
|   009|      Tom Tan| 33|58000.0|11600.0|        1|      two|
|   010|     Lisa Lee| 27|47000.0| 9400.0|        1|      two|
|   011|   David Park| 38|65000.0|13000.0|        1|      two|
|   012|   Susan Chen| 31|54000.0|10800.0|        1|      two|
|   013|    Brian Kim| 45|75000.0|15000.0|        1|   

In [17]:
emp_2 = emp_new_cols.withColumnRenamed(
    "ColumnTwo", 'Column Two'
)

In [19]:
emp_2.show()

+-----------+-------------+---+-------+-------+---------+----------+
|employee_id|         name|age| salary|    tax|ColumnOne|Volumn Two|
+-----------+-------------+---+-------+-------+---------+----------+
|        001|     John Doe| 30|50000.0|10000.0|        1|       two|
|        002|   Jane Smith| 25|45000.0| 9000.0|        1|       two|
|        003|    Bob Brown| 35|55000.0|11000.0|        1|       two|
|        004|    Alice Lee| 28|48000.0| 9600.0|        1|       two|
|        005|    Jack Chan| 40|60000.0|12000.0|        1|       two|
|        006|    Jill Wong| 32|52000.0|10400.0|        1|       two|
|        007|James Johnson| 42|70000.0|14000.0|        1|       two|
|        008|     Kate Kim| 29|51000.0|10200.0|        1|       two|
|        009|      Tom Tan| 33|58000.0|11600.0|        1|       two|
|        010|     Lisa Lee| 27|47000.0| 9400.0|        1|       two|
|        011|   David Park| 38|65000.0|13000.0|        1|       two|
|        012|   Susan Chen| 31|540

In [20]:
emp_dropped = emp_new_cols.drop('ColumnTwo')

In [21]:
emp_dropped.show()

+-----------+-------------+---+-------+-------+---------+
|employee_id|         name|age| salary|    tax|ColumnOne|
+-----------+-------------+---+-------+-------+---------+
|        001|     John Doe| 30|50000.0|10000.0|        1|
|        002|   Jane Smith| 25|45000.0| 9000.0|        1|
|        003|    Bob Brown| 35|55000.0|11000.0|        1|
|        004|    Alice Lee| 28|48000.0| 9600.0|        1|
|        005|    Jack Chan| 40|60000.0|12000.0|        1|
|        006|    Jill Wong| 32|52000.0|10400.0|        1|
|        007|James Johnson| 42|70000.0|14000.0|        1|
|        008|     Kate Kim| 29|51000.0|10200.0|        1|
|        009|      Tom Tan| 33|58000.0|11600.0|        1|
|        010|     Lisa Lee| 27|47000.0| 9400.0|        1|
|        011|   David Park| 38|65000.0|13000.0|        1|
|        012|   Susan Chen| 31|54000.0|10800.0|        1|
|        013|    Brian Kim| 45|75000.0|15000.0|        1|
|        014|    Emily Lee| 26|46000.0| 9200.0|        1|
|        015| 

In [22]:
emp_dropped = emp_new_cols.drop('ColumnOne', 'ColumnTwo')

In [24]:
emp_dropped.show(5)

+-----------+----------+---+-------+-------+
|employee_id|      name|age| salary|    tax|
+-----------+----------+---+-------+-------+
|        001|  John Doe| 30|50000.0|10000.0|
|        002|Jane Smith| 25|45000.0| 9000.0|
|        003| Bob Brown| 35|55000.0|11000.0|
|        004| Alice Lee| 28|48000.0| 9600.0|
|        005| Jack Chan| 40|60000.0|12000.0|
+-----------+----------+---+-------+-------+
only showing top 5 rows



In [25]:
emp_filtered = emp_dropped.where('tax > 10000')

In [26]:
emp_filtered.show(5)

+-----------+-------------+---+-------+-------+
|employee_id|         name|age| salary|    tax|
+-----------+-------------+---+-------+-------+
|        003|    Bob Brown| 35|55000.0|11000.0|
|        005|    Jack Chan| 40|60000.0|12000.0|
|        006|    Jill Wong| 32|52000.0|10400.0|
|        007|James Johnson| 42|70000.0|14000.0|
|        008|     Kate Kim| 29|51000.0|10200.0|
+-----------+-------------+---+-------+-------+
only showing top 5 rows



In [27]:
emp_limit = emp_filtered.limit(5)

In [28]:
emp_limit.show(5)

+-----------+-------------+---+-------+-------+
|employee_id|         name|age| salary|    tax|
+-----------+-------------+---+-------+-------+
|        003|    Bob Brown| 35|55000.0|11000.0|
|        005|    Jack Chan| 40|60000.0|12000.0|
|        006|    Jill Wong| 32|52000.0|10400.0|
|        007|James Johnson| 42|70000.0|14000.0|
|        008|     Kate Kim| 29|51000.0|10200.0|
+-----------+-------------+---+-------+-------+



In [29]:
columns = {
    'tax': F.col('salary') * 0.2,
    'one_number': F.lit(1),
    'two_number': F.lit('two')
}

In [30]:
emp_final = emp_filtered.withColumns(
    columns
)

In [31]:
emp_final.show(5)

+-----------+-------------+---+-------+-------+----------+----------+
|employee_id|         name|age| salary|    tax|one_number|two_number|
+-----------+-------------+---+-------+-------+----------+----------+
|        003|    Bob Brown| 35|55000.0|11000.0|         1|       two|
|        005|    Jack Chan| 40|60000.0|12000.0|         1|       two|
|        006|    Jill Wong| 32|52000.0|10400.0|         1|       two|
|        007|James Johnson| 42|70000.0|14000.0|         1|       two|
|        008|     Kate Kim| 29|51000.0|10200.0|         1|       two|
+-----------+-------------+---+-------+-------+----------+----------+
only showing top 5 rows

