In [0]:
from pyspark.sql import SparkSession
from pyspark.sql.functions import *
from pyspark import SparkContext, SparkConf

spark = SparkSession.builder.appName("Union").master("local[*]").getOrCreate()


In [0]:
spark

In [0]:
data_1 = [
    ["000", "107", "Emily Lee", 26, "", 46000, "2019-01-01"],
    ["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"],
    ]

data_2 = [
    ["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"],
]

emp_schema = "employee_id string, department_id string, name string, age string, gender string, salary string, hire_date string"

In [0]:
emp_data_1 = spark.createDataFrame(data_1,schema=emp_schema)
emp_data_2 = spark.createDataFrame(data_2,schema= emp_schema)

emp_data_1.show()
emp_data_2.show()

+-----------+-------------+-------------+---+------+------+----------+
|employee_id|department_id|         name|age|gender|salary| hire_date|
+-----------+-------------+-------------+---+------+------+----------+
|        000|          107|    Emily Lee| 26|      | 46000|2019-01-01|
|        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|
+-----

In [0]:
emp_data_1.printSchema()
emp_data_2.printSchema()

root
 |-- employee_id: string (nullable = true)
 |-- department_id: string (nullable = true)
 |-- name: string (nullable = true)
 |-- age: string (nullable = true)
 |-- gender: string (nullable = true)
 |-- salary: string (nullable = true)
 |-- hire_date: string (nullable = true)

root
 |-- employee_id: string (nullable = true)
 |-- department_id: string (nullable = true)
 |-- name: string (nullable = true)
 |-- age: string (nullable = true)
 |-- gender: string (nullable = true)
 |-- salary: string (nullable = true)
 |-- hire_date: string (nullable = true)



# UNION
## Union used to combine both tables and that both the tables should have same schema, Columns should be same and column sequences should same

## Incase of Union only distinct values will be given no duplicates, but incase of unionall all the values from both the tables are given.

In [0]:
emp = emp_data_1.union(emp_data_2)
emp.show()
emp_unionAll = emp_data_1.unionAll(emp_data_2).show()

+-----------+-------------+-------------+---+------+------+----------+
|employee_id|department_id|         name|age|gender|salary| hire_date|
+-----------+-------------+-------------+---+------+------+----------+
|        000|          107|    Emily Lee| 26|      | 46000|2019-01-01|
|        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|
|     

## Sort the employees based on salary desc

In [0]:
emp_salary_desc = emp.orderBy(col("salary").desc())
emp_salary_desc.show()

+-----------+-------------+-------------+---+------+------+----------+
|employee_id|department_id|         name|age|gender|salary| hire_date|
+-----------+-------------+-------------+---+------+------+----------+
|        013|          106|    Brian Kim| 45|  Male| 75000|2011-07-01|
|        007|          101|James Johnson| 42|  Male| 70000|2012-03-15|
|        011|          104|   David Park| 38|  Male| 65000|2015-11-01|
|        015|          106|  Michael Lee| 37|  Male| 63000|2014-09-30|
|        019|          103|  Steven Chen| 36|  Male| 62000|2015-08-01|
|        005|          103|    Jack Chan| 40|  Male| 60000|2013-04-01|
|        009|          103|      Tom Tan| 33|  Male| 58000|2016-06-01|
|        017|          105|  George Wang| 34|  Male| 57000|2016-03-15|
|        003|          102|    Bob Brown| 35|  Male| 55000|2014-05-01|
|        012|          105|   Susan Chen| 31|Female| 54000|2017-02-15|
|        020|          102|    Grace Kim| 32|Female| 53000|2018-11-01|
|     

In [0]:
emp_salary_asc = emp.orderBy(col("salary").asc()).show()

+-----------+-------------+-------------+---+------+------+----------+
|employee_id|department_id|         name|age|gender|salary| hire_date|
+-----------+-------------+-------------+---+------+------+----------+
|        002|          101|   Jane Smith| 25|Female| 45000|2016-02-15|
|        014|          107|    Emily Lee| 26|Female| 46000|2019-01-01|
|        000|          107|    Emily Lee| 26|      | 46000|2019-01-01|
|        010|          104|     Lisa Lee| 27|Female| 47000|2018-08-01|
|        004|          102|    Alice Lee| 28|Female| 48000|2017-09-30|
|        016|          107|  Kelly Zhang| 30|Female| 49000|2018-04-01|
|        018|          104|    Nancy Liu| 29|Female| 50000|2017-06-01|
|        001|          101|     John Doe| 30|  Male| 50000|2015-01-01|
|        008|          102|     Kate Kim| 29|Female| 51000|2019-10-01|
|        006|          103|    Jill Wong| 32|Female| 52000|2018-07-01|
|        020|          102|    Grace Kim| 32|Female| 53000|2018-11-01|
|     

# Aggregating the Data based on department i'd and count no.of employees in that department

In [0]:
emp_aggr = emp.groupBy("department_id").agg(count("employee_id").alias("total_department_count"))
emp_aggr.show()

+-------------+----------------------+
|department_id|total_department_count|
+-------------+----------------------+
|          107|                     3|
|          101|                     3|
|          102|                     4|
|          103|                     4|
|          104|                     3|
|          105|                     2|
|          106|                     2|
+-------------+----------------------+



## Aggregate the data by grouping based on department i'd and sum the salaries in department wise.

In [0]:
emp_salary_group = emp.groupBy("department_id").agg(sum("salary").alias("sum_of_salaries")).orderBy(col("sum_of_salaries").asc())
emp_salary_group.show()

+-------------+---------------+
|department_id|sum_of_salaries|
+-------------+---------------+
|          105|       111000.0|
|          106|       138000.0|
|          107|       141000.0|
|          104|       162000.0|
|          101|       165000.0|
|          102|       207000.0|
|          103|       232000.0|
+-------------+---------------+



## Aggregate the data using grouping on department i'd and using having clause find the departments whose having avg.salary > 50000 

In [0]:
emp_having = emp.groupBy("department_id").agg(avg("salary").alias("avg_salary")).where("avg_salary > 50000")
emp_having.show()

emp_having2 = emp.groupBy("department_id").agg(avg("salary").alias("avg_salary")).filter((col("avg_salary") > 50000))
emp_having2.show() 

+-------------+----------+
|department_id|avg_salary|
+-------------+----------+
|          101|   55000.0|
|          102|   51750.0|
|          103|   58000.0|
|          104|   54000.0|
|          105|   55500.0|
|          106|   69000.0|
+-------------+----------+

+-------------+----------+
|department_id|avg_salary|
+-------------+----------+
|          101|   55000.0|
|          102|   51750.0|
|          103|   58000.0|
|          104|   54000.0|
|          105|   55500.0|
|          106|   69000.0|
+-------------+----------+



# Using unionByName we can combine two tables even the columns are not in sequential order, unionByName matches with column name and tries to combine.

In [0]:
emp.count()
# spark by default counts the no.of records in a data frame

Out[44]: 21