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

In [2]:
from pyspark.conf import SparkConf
config = SparkConf()
config.setMaster("local").setAppName("DFJoin")


from pyspark.sql import SparkSession
spark = SparkSession.builder.config(conf=config).getOrCreate()

22/06/03 14:35:20 WARN Utils: Your hostname, Hocines-MacBook-Pro.local resolves to a loopback address: 127.0.0.1; using 10.0.0.164 instead (on interface en0)
22/06/03 14:35:20 WARN Utils: Set SPARK_LOCAL_IP if you need to bind to another address
Using Spark's default log4j profile: org/apache/spark/log4j-defaults.properties
Setting default log level to "WARN".
To adjust logging level use sc.setLogLevel(newLevel). For SparkR, use setLogLevel(newLevel).
22/06/03 14:35:20 WARN NativeCodeLoader: Unable to load native-hadoop library for your platform... using builtin-java classes where applicable
22/06/03 14:35:21 WARN Utils: Service 'SparkUI' could not bind on port 4040. Attempting port 4041.


In [3]:
products = [ 
          # (product_id, product_name, brand_id)  
         (1, 'iPhone', 100),
         (2, 'Galaxy', 200),
         (3, 'Redme', 300), # orphan record, no matching brand
         (4, 'Pixel', 400),
]

brands = [
    #(brand_id, brand_name)
    (100, "Apple"),
    (200, "Samsung"),
    (400, "Google"),
    (500, "Sony"), # no matching products
]
 
productDf = spark.createDataFrame(data=products, schema=["product_id", "product_name", "brand_id"])
brandDf = spark.createDataFrame(data=brands, schema=["brand_id", "brand_name"])
productDf.show()
brandDf.show()


productDf.createOrReplaceTempView("products")
brandDf.createOrReplaceTempView("brands")


[Stage 0:>                                                          (0 + 1) / 1]

+----------+------------+--------+
|product_id|product_name|brand_id|
+----------+------------+--------+
|         1|      iPhone|     100|
|         2|      Galaxy|     200|
|         3|       Redme|     300|
|         4|       Pixel|     400|
+----------+------------+--------+

+--------+----------+
|brand_id|brand_name|
+--------+----------+
|     100|     Apple|
|     200|   Samsung|
|     400|    Google|
|     500|      Sony|
+--------+----------+



                                                                                

In [6]:
spark.sql("""
SELECT products.*, brands.brand_name FROM products
INNER JOIN brands ON products.brand_id = brands.brand_id
"""). show()

+----------+------------+--------+----------+
|product_id|product_name|brand_id|brand_name|
+----------+------------+--------+----------+
|         1|      iPhone|     100|     Apple|
|         2|      Galaxy|     200|   Samsung|
|         4|       Pixel|     400|    Google|
+----------+------------+--------+----------+



In [4]:
# Inner Join
# productDf is left
# brandDf is right
# select/pick only matching record, discord if no matches found
productDf.join(brandDf, productDf["brand_id"] ==  brandDf["brand_id"], "inner").show()

                                                                                

+----------+------------+--------+--------+----------+
|product_id|product_name|brand_id|brand_id|brand_name|
+----------+------------+--------+--------+----------+
|         1|      iPhone|     100|     100|     Apple|
|         2|      Galaxy|     200|     200|   Samsung|
|         4|       Pixel|     400|     400|    Google|
+----------+------------+--------+--------+----------+



In [7]:
spark.sql("""
SELECT products.*, brands.brand_name FROM products
FULL OUTER JOIN brands ON products.brand_id = brands.brand_id
"""). show()

+----------+------------+--------+----------+
|product_id|product_name|brand_id|brand_name|
+----------+------------+--------+----------+
|         1|      iPhone|     100|     Apple|
|         2|      Galaxy|     200|   Samsung|
|         3|       Redme|     300|      null|
|         4|       Pixel|     400|    Google|
|      null|        null|    null|      Sony|
+----------+------------+--------+----------+



In [5]:
# Outer Join, Full Outer Outer, [Left outer + Right outer]
# pick all records from left dataframe, and also right dataframe
# if no matches found, it fills null data for not matched records
productDf.join(brandDf, productDf["brand_id"] ==  brandDf["brand_id"], "outer").show()

                                                                                

+----------+------------+--------+--------+----------+
|product_id|product_name|brand_id|brand_id|brand_name|
+----------+------------+--------+--------+----------+
|      null|        null|    null|     500|      Sony|
|         1|      iPhone|     100|     100|     Apple|
|         2|      Galaxy|     200|     200|   Samsung|
|         4|       Pixel|     400|     400|    Google|
|         3|       Redme|     300|    null|      null|
+----------+------------+--------+--------+----------+



                                                                                

In [8]:
spark.sql("""
SELECT products.*, brands.brand_name FROM products
LEFT OUTER JOIN brands ON products.brand_id = brands.brand_id
"""). show()

+----------+------------+--------+----------+
|product_id|product_name|brand_id|brand_name|
+----------+------------+--------+----------+
|         1|      iPhone|     100|     Apple|
|         2|      Galaxy|     200|   Samsung|
|         3|       Redme|     300|      null|
|         4|       Pixel|     400|    Google|
+----------+------------+--------+----------+



In [6]:
# Left, Left Outer join 
# picks all records from left, if no matches found, it fills null for right data
productDf.join(brandDf, productDf["brand_id"] ==  brandDf["brand_id"], "leftouter").show()

                                                                                

+----------+------------+--------+--------+----------+
|product_id|product_name|brand_id|brand_id|brand_name|
+----------+------------+--------+--------+----------+
|         1|      iPhone|     100|     100|     Apple|
|         2|      Galaxy|     200|     200|   Samsung|
|         4|       Pixel|     400|     400|    Google|
|         3|       Redme|     300|    null|      null|
+----------+------------+--------+--------+----------+



In [9]:
spark.sql("""
SELECT products.*, brands.brand_name FROM products
RIGHT OUTER JOIN brands ON products.brand_id = brands.brand_id
"""). show()

+----------+------------+--------+----------+
|product_id|product_name|brand_id|brand_name|
+----------+------------+--------+----------+
|         1|      iPhone|     100|     Apple|
|         2|      Galaxy|     200|   Samsung|
|         4|       Pixel|     400|    Google|
|      null|        null|    null|      Sony|
+----------+------------+--------+----------+



In [7]:
# Right, Right outer Join
# picks all the records from right, if no matches found, fills left data with null
productDf.join(brandDf, productDf["brand_id"] ==  brandDf["brand_id"], "rightouter").show()

                                                                                

+----------+------------+--------+--------+----------+
|product_id|product_name|brand_id|brand_id|brand_name|
+----------+------------+--------+--------+----------+
|      null|        null|    null|     500|      Sony|
|         1|      iPhone|     100|     100|     Apple|
|         2|      Galaxy|     200|     200|   Samsung|
|         4|       Pixel|     400|     400|    Google|
+----------+------------+--------+--------+----------+



In [9]:
store = [
    #(store_id, store_name)
    (1000, "Poorvika"),
    (2000, "Sangeetha"),
    (4000, "Amazon"),
    (5000, "FlipKart"), 
]
 
storeDf = spark.createDataFrame(data=store, schema=["store_id", "store_name"])
storeDf.show()

+--------+----------+
|store_id|store_name|
+--------+----------+
|    1000|  Poorvika|
|    2000| Sangeetha|
|    4000|    Amazon|
|    5000|  FlipKart|
+--------+----------+



In [10]:
# cartesian , take row from left side, pair with all from right side
productDf.crossJoin(storeDf).show()


+----------+------------+--------+--------+----------+
|product_id|product_name|brand_id|store_id|store_name|
+----------+------------+--------+--------+----------+
|         1|      iPhone|     100|    1000|  Poorvika|
|         1|      iPhone|     100|    2000| Sangeetha|
|         1|      iPhone|     100|    4000|    Amazon|
|         1|      iPhone|     100|    5000|  FlipKart|
|         2|      Galaxy|     200|    1000|  Poorvika|
|         2|      Galaxy|     200|    2000| Sangeetha|
|         2|      Galaxy|     200|    4000|    Amazon|
|         2|      Galaxy|     200|    5000|  FlipKart|
|         3|       Redme|     300|    1000|  Poorvika|
|         3|       Redme|     300|    2000| Sangeetha|
|         3|       Redme|     300|    4000|    Amazon|
|         3|       Redme|     300|    5000|  FlipKart|
|         4|       Pixel|     400|    1000|  Poorvika|
|         4|       Pixel|     400|    2000| Sangeetha|
|         4|       Pixel|     400|    4000|    Amazon|
|         