# Spark API

In [1]:
import pyspark

spark = pyspark.sql.SparkSession.builder.getOrCreate()

In [2]:
spark

In [3]:
#for demonstration, we'll create a spark dataframe from a pandas dataframe
import numpy as np
import pandas as pd
import pydataset

In [4]:
#import tips from pydataset
tips = pydataset.data('tips')

#turn pandas df into a spark df
df = spark.createDataFrame(tips)

#look at the data
df.show()

+----------+----+------+------+---+------+----+
|total_bill| tip|   sex|smoker|day|  time|size|
+----------+----+------+------+---+------+----+
|     16.99|1.01|Female|    No|Sun|Dinner|   2|
|     10.34|1.66|  Male|    No|Sun|Dinner|   3|
|     21.01| 3.5|  Male|    No|Sun|Dinner|   3|
|     23.68|3.31|  Male|    No|Sun|Dinner|   2|
|     24.59|3.61|Female|    No|Sun|Dinner|   4|
|     25.29|4.71|  Male|    No|Sun|Dinner|   4|
|      8.77| 2.0|  Male|    No|Sun|Dinner|   2|
|     26.88|3.12|  Male|    No|Sun|Dinner|   4|
|     15.04|1.96|  Male|    No|Sun|Dinner|   2|
|     14.78|3.23|  Male|    No|Sun|Dinner|   2|
|     10.27|1.71|  Male|    No|Sun|Dinner|   2|
|     35.26| 5.0|Female|    No|Sun|Dinner|   4|
|     15.42|1.57|  Male|    No|Sun|Dinner|   2|
|     18.43| 3.0|  Male|    No|Sun|Dinner|   4|
|     14.83|3.02|Female|    No|Sun|Dinner|   2|
|     21.58|3.92|  Male|    No|Sun|Dinner|   2|
|     10.33|1.67|Female|    No|Sun|Dinner|   3|
|     16.29|3.71|  Male|    No|Sun|Dinne

<hr style="border:2px solid black"> </hr>

## DataFrame Basics

In [5]:
#spark is lazy
#if will only show you the column names and size
df

DataFrame[total_bill: double, tip: double, sex: string, smoker: string, day: string, time: string, size: bigint]

In [6]:
#this is how you actually look at the contents
df.show(5)

+----------+----+------+------+---+------+----+
|total_bill| tip|   sex|smoker|day|  time|size|
+----------+----+------+------+---+------+----+
|     16.99|1.01|Female|    No|Sun|Dinner|   2|
|     10.34|1.66|  Male|    No|Sun|Dinner|   3|
|     21.01| 3.5|  Male|    No|Sun|Dinner|   3|
|     23.68|3.31|  Male|    No|Sun|Dinner|   2|
|     24.59|3.61|Female|    No|Sun|Dinner|   4|
+----------+----+------+------+---+------+----+
only showing top 5 rows



In [8]:
#this is how you transpose
df.show(5, vertical=True)

-RECORD 0------------
 total_bill | 16.99  
 tip        | 1.01   
 sex        | Female 
 smoker     | No     
 day        | Sun    
 time       | Dinner 
 size       | 2      
-RECORD 1------------
 total_bill | 10.34  
 tip        | 1.66   
 sex        | Male   
 smoker     | No     
 day        | Sun    
 time       | Dinner 
 size       | 3      
-RECORD 2------------
 total_bill | 21.01  
 tip        | 3.5    
 sex        | Male   
 smoker     | No     
 day        | Sun    
 time       | Dinner 
 size       | 3      
-RECORD 3------------
 total_bill | 23.68  
 tip        | 3.31   
 sex        | Male   
 smoker     | No     
 day        | Sun    
 time       | Dinner 
 size       | 2      
-RECORD 4------------
 total_bill | 24.59  
 tip        | 3.61   
 sex        | Female 
 smoker     | No     
 day        | Sun    
 time       | Dinner 
 size       | 4      
only showing top 5 rows



In [9]:
# Don't do this!
# just use .show to view df contents
df2 = df.show(10)

+----------+----+------+------+---+------+----+
|total_bill| tip|   sex|smoker|day|  time|size|
+----------+----+------+------+---+------+----+
|     16.99|1.01|Female|    No|Sun|Dinner|   2|
|     10.34|1.66|  Male|    No|Sun|Dinner|   3|
|     21.01| 3.5|  Male|    No|Sun|Dinner|   3|
|     23.68|3.31|  Male|    No|Sun|Dinner|   2|
|     24.59|3.61|Female|    No|Sun|Dinner|   4|
|     25.29|4.71|  Male|    No|Sun|Dinner|   4|
|      8.77| 2.0|  Male|    No|Sun|Dinner|   2|
|     26.88|3.12|  Male|    No|Sun|Dinner|   4|
|     15.04|1.96|  Male|    No|Sun|Dinner|   2|
|     14.78|3.23|  Male|    No|Sun|Dinner|   2|
+----------+----+------+------+---+------+----+
only showing top 10 rows



In [10]:
#df.show() == is a print statement
#can not save a variable to .show()
type(df2)

NoneType

In [11]:
#gives back a list of row objects
df.head(5)

[Row(total_bill=16.99, tip=1.01, sex='Female', smoker='No', day='Sun', time='Dinner', size=2),
 Row(total_bill=10.34, tip=1.66, sex='Male', smoker='No', day='Sun', time='Dinner', size=3),
 Row(total_bill=21.01, tip=3.5, sex='Male', smoker='No', day='Sun', time='Dinner', size=3),
 Row(total_bill=23.68, tip=3.31, sex='Male', smoker='No', day='Sun', time='Dinner', size=2),
 Row(total_bill=24.59, tip=3.61, sex='Female', smoker='No', day='Sun', time='Dinner', size=4)]

In [12]:
#grab first item from the list
df.head(5)[0]

Row(total_bill=16.99, tip=1.01, sex='Female', smoker='No', day='Sun', time='Dinner', size=2)

In [13]:
#grab just time from that row
df.head(5)[0].time

'Dinner'

<hr style="border:2px solid black"> </hr>

## Selecting Columns

In [15]:
#filter down to specified columns
df.select('total_bill', 'tip', 'size', 'day').show(5)

+----------+----+----+---+
|total_bill| tip|size|day|
+----------+----+----+---+
|     16.99|1.01|   2|Sun|
|     10.34|1.66|   3|Sun|
|     21.01| 3.5|   3|Sun|
|     23.68|3.31|   2|Sun|
|     24.59|3.61|   4|Sun|
+----------+----+----+---+
only showing top 5 rows



In [16]:
#returns everything
df.select('*')

DataFrame[total_bill: double, tip: double, sex: string, smoker: string, day: string, time: string, size: bigint]

In [17]:
#perform division in table
df.select(df.tip / df.total_bill).show(5)

+-------------------+
| (tip / total_bill)|
+-------------------+
|0.05944673337257211|
|0.16054158607350097|
|0.16658733936220846|
| 0.1397804054054054|
|0.14680764538430255|
+-------------------+
only showing top 5 rows



In [18]:
#assign to variable
col = df.tip / df.total_bill
col

Column<'(tip / total_bill)'>

In [22]:
#create new column
df.select('*', col).show(5)

+----------+----+------+------+---+------+----+-------------------+
|total_bill| tip|   sex|smoker|day|  time|size| (tip / total_bill)|
+----------+----+------+------+---+------+----+-------------------+
|     16.99|1.01|Female|    No|Sun|Dinner|   2|0.05944673337257211|
|     10.34|1.66|  Male|    No|Sun|Dinner|   3|0.16054158607350097|
|     21.01| 3.5|  Male|    No|Sun|Dinner|   3|0.16658733936220846|
|     23.68|3.31|  Male|    No|Sun|Dinner|   2| 0.1397804054054054|
|     24.59|3.61|Female|    No|Sun|Dinner|   4|0.14680764538430255|
+----------+----+------+------+---+------+----+-------------------+
only showing top 5 rows



In [23]:
#rename new column using .alias
df.select('*', col.alias('tip_pct')).show(5)

+----------+----+------+------+---+------+----+-------------------+
|total_bill| tip|   sex|smoker|day|  time|size|            tip_pct|
+----------+----+------+------+---+------+----+-------------------+
|     16.99|1.01|Female|    No|Sun|Dinner|   2|0.05944673337257211|
|     10.34|1.66|  Male|    No|Sun|Dinner|   3|0.16054158607350097|
|     21.01| 3.5|  Male|    No|Sun|Dinner|   3|0.16658733936220846|
|     23.68|3.31|  Male|    No|Sun|Dinner|   2| 0.1397804054054054|
|     24.59|3.61|Female|    No|Sun|Dinner|   4|0.14680764538430255|
+----------+----+------+------+---+------+----+-------------------+
only showing top 5 rows



In [24]:
#assign to new variable
df_with_tip_pct = df.select('*', col.alias('tip_pct'))

In [25]:
#look at new df with added column
df_with_tip_pct.show(5)

+----------+----+------+------+---+------+----+-------------------+
|total_bill| tip|   sex|smoker|day|  time|size|            tip_pct|
+----------+----+------+------+---+------+----+-------------------+
|     16.99|1.01|Female|    No|Sun|Dinner|   2|0.05944673337257211|
|     10.34|1.66|  Male|    No|Sun|Dinner|   3|0.16054158607350097|
|     21.01| 3.5|  Male|    No|Sun|Dinner|   3|0.16658733936220846|
|     23.68|3.31|  Male|    No|Sun|Dinner|   2| 0.1397804054054054|
|     24.59|3.61|Female|    No|Sun|Dinner|   4|0.14680764538430255|
+----------+----+------+------+---+------+----+-------------------+
only showing top 5 rows



<hr style="border:2px solid black"> </hr>

## Selecting w/ Built In Functions

In [26]:
#import built in functions
from pyspark.sql.functions import sum, mean, concat, lit, regexp_extract, regexp_replace, when

In [27]:
#average tip, sum of all total bills
df.select(mean(df.tip), sum(df.total_bill)).show()

+------------------+-----------------+
|          avg(tip)|  sum(total_bill)|
+------------------+-----------------+
|2.9982786885245907|4827.769999999999|
+------------------+-----------------+



In [28]:
#add weekday and time of service with space in between
#can only concat string value columns
df.select(concat('day', lit(' '), 'time')).show(5)

+--------------------+
|concat(day,  , time)|
+--------------------+
|          Sun Dinner|
|          Sun Dinner|
|          Sun Dinner|
|          Sun Dinner|
|          Sun Dinner|
+--------------------+
only showing top 5 rows



In [30]:
#changing time to int value
df.select(df.time.cast('int')).show(5)

#gives nulls because it cannot change a string to an int

+----+
|time|
+----+
|null|
|null|
|null|
|null|
|null|
+----+
only showing top 5 rows



In [31]:
#change size to string value, was int
df.select(df.size.cast('string')).show(5)

+----+
|size|
+----+
|   2|
|   3|
|   3|
|   2|
|   4|
+----+
only showing top 5 rows



- **.cast() = transformation**
- **.show() = action**
- **.count() = action**

In [32]:
#reassign existing df to show percent column
df = df.select(
    '*',
    (df.tip / df.total_bill).alias('tip_pct')
)

In [33]:
df.show()

+----------+----+------+------+---+------+----+-------------------+
|total_bill| tip|   sex|smoker|day|  time|size|            tip_pct|
+----------+----+------+------+---+------+----+-------------------+
|     16.99|1.01|Female|    No|Sun|Dinner|   2|0.05944673337257211|
|     10.34|1.66|  Male|    No|Sun|Dinner|   3|0.16054158607350097|
|     21.01| 3.5|  Male|    No|Sun|Dinner|   3|0.16658733936220846|
|     23.68|3.31|  Male|    No|Sun|Dinner|   2| 0.1397804054054054|
|     24.59|3.61|Female|    No|Sun|Dinner|   4|0.14680764538430255|
|     25.29|4.71|  Male|    No|Sun|Dinner|   4|0.18623962040332148|
|      8.77| 2.0|  Male|    No|Sun|Dinner|   2|0.22805017103762829|
|     26.88|3.12|  Male|    No|Sun|Dinner|   4|0.11607142857142858|
|     15.04|1.96|  Male|    No|Sun|Dinner|   2|0.13031914893617022|
|     14.78|3.23|  Male|    No|Sun|Dinner|   2| 0.2185385656292287|
|     10.27|1.71|  Male|    No|Sun|Dinner|   2| 0.1665043816942551|
|     35.26| 5.0|Female|    No|Sun|Dinner|   4|0

<hr style="border:2px solid black"> </hr>

## When/ Otherwise

Similar to:
- if/else in excel
- np.where in Numpy


In [34]:
df.select(
    'tip_pct',
).show(5)

+-------------------+
|            tip_pct|
+-------------------+
|0.05944673337257211|
|0.16054158607350097|
|0.16658733936220846|
| 0.1397804054054054|
|0.14680764538430255|
+-------------------+
only showing top 5 rows



In [35]:
#when tip is greater then 20% = good tip
when(df.tip_pct > .2, 'good tip')

Column<'CASE WHEN (tip_pct > 0.2) THEN good tip END'>

In [37]:
#nulls = do not satisfy greater than 20%
df.select(
    'tip_pct',
    (when(df.tip_pct > .2, 'good tip')
     .alias('tip_desc'))
).show(25)

+-------------------+--------+
|            tip_pct|tip_desc|
+-------------------+--------+
|0.05944673337257211|    null|
|0.16054158607350097|    null|
|0.16658733936220846|    null|
| 0.1397804054054054|    null|
|0.14680764538430255|    null|
|0.18623962040332148|    null|
|0.22805017103762829|good tip|
|0.11607142857142858|    null|
|0.13031914893617022|    null|
| 0.2185385656292287|good tip|
| 0.1665043816942551|    null|
|0.14180374361883155|    null|
|0.10181582360570687|    null|
|0.16277807921866522|    null|
|0.20364126770060686|good tip|
|0.18164967562557924|    null|
| 0.1616650532429816|    null|
|0.22774708410067526|good tip|
|0.20624631703005306|good tip|
|0.16222760290556903|    null|
|0.22767857142857142|good tip|
|0.13553474618038444|    null|
|0.14140773620798985|    null|
|0.19228817858954844|    null|
|0.16044399596367306|    null|
+-------------------+--------+
only showing top 25 rows



In [38]:
#use .otherwise() to get rid of nulls
#tip greater then 20% = good tip
#other wise = not good tip
#rename column
df.select(
    'tip_pct',
    (when(df.tip_pct > .2, 'good tip')
     .otherwise('not good tip')
     .alias('tip_desc'))
).show(25)

+-------------------+------------+
|            tip_pct|    tip_desc|
+-------------------+------------+
|0.05944673337257211|not good tip|
|0.16054158607350097|not good tip|
|0.16658733936220846|not good tip|
| 0.1397804054054054|not good tip|
|0.14680764538430255|not good tip|
|0.18623962040332148|not good tip|
|0.22805017103762829|    good tip|
|0.11607142857142858|not good tip|
|0.13031914893617022|not good tip|
| 0.2185385656292287|    good tip|
| 0.1665043816942551|not good tip|
|0.14180374361883155|not good tip|
|0.10181582360570687|not good tip|
|0.16277807921866522|not good tip|
|0.20364126770060686|    good tip|
|0.18164967562557924|not good tip|
| 0.1616650532429816|not good tip|
|0.22774708410067526|    good tip|
|0.20624631703005306|    good tip|
|0.16222760290556903|not good tip|
|0.22767857142857142|    good tip|
|0.13553474618038444|not good tip|
|0.14140773620798985|not good tip|
|0.19228817858954844|not good tip|
|0.16044399596367306|not good tip|
+-------------------

<hr style="border:2px solid black"> </hr>

## Regex

- extract = pull piece out of regex
- replace = replace a piece of a sting

In [39]:
#using 'time' column
#pulls out first letter and puts it into a new column called 'first_letter'
#replaces any vowel with X

df.select(
    'time',
    regexp_extract('time', r'(\w).*', 1).alias('first_letter'),
    regexp_replace('time', r'[aeiou]', 'X')
).show(5)

+------+------------+-----------------------------------+
|  time|first_letter|regexp_replace(time, [aeiou], X, 1)|
+------+------------+-----------------------------------+
|Dinner|           D|                             DXnnXr|
|Dinner|           D|                             DXnnXr|
|Dinner|           D|                             DXnnXr|
|Dinner|           D|                             DXnnXr|
|Dinner|           D|                             DXnnXr|
+------+------------+-----------------------------------+
only showing top 5 rows



In [42]:
#using 'sex' only
#extract first letter
#replace vowels with -
df.select(
    'sex',
    regexp_extract('sex', r'(\w).*', 1).alias('first_letter'),
    regexp_replace('sex', r'[aeiou]', '-')
).show(5)

+------+------------+----------------------------------+
|   sex|first_letter|regexp_replace(sex, [aeiou], -, 1)|
+------+------------+----------------------------------+
|Female|           F|                            F-m-l-|
|  Male|           M|                              M-l-|
|  Male|           M|                              M-l-|
|  Male|           M|                              M-l-|
|Female|           F|                            F-m-l-|
+------+------------+----------------------------------+
only showing top 5 rows



<hr style="border:2px solid black"> </hr>

## Transforming Rows

In [43]:
df.show()

+----------+----+------+------+---+------+----+-------------------+
|total_bill| tip|   sex|smoker|day|  time|size|            tip_pct|
+----------+----+------+------+---+------+----+-------------------+
|     16.99|1.01|Female|    No|Sun|Dinner|   2|0.05944673337257211|
|     10.34|1.66|  Male|    No|Sun|Dinner|   3|0.16054158607350097|
|     21.01| 3.5|  Male|    No|Sun|Dinner|   3|0.16658733936220846|
|     23.68|3.31|  Male|    No|Sun|Dinner|   2| 0.1397804054054054|
|     24.59|3.61|Female|    No|Sun|Dinner|   4|0.14680764538430255|
|     25.29|4.71|  Male|    No|Sun|Dinner|   4|0.18623962040332148|
|      8.77| 2.0|  Male|    No|Sun|Dinner|   2|0.22805017103762829|
|     26.88|3.12|  Male|    No|Sun|Dinner|   4|0.11607142857142858|
|     15.04|1.96|  Male|    No|Sun|Dinner|   2|0.13031914893617022|
|     14.78|3.23|  Male|    No|Sun|Dinner|   2| 0.2185385656292287|
|     10.27|1.71|  Male|    No|Sun|Dinner|   2| 0.1665043816942551|
|     35.26| 5.0|Female|    No|Sun|Dinner|   4|0

<hr style="border:2px solid black"> </hr>

## Sorting
- .orderBy
- .sort
<br>

- sorting is VERY expensive
- use as LAST STEP or on subset only
- can specify if you want Nulls first or last

In [44]:
#order by total_bill - smallest to largest
df.orderBy(df.total_bill).show()

+----------+----+------+------+----+------+----+-------------------+
|total_bill| tip|   sex|smoker| day|  time|size|            tip_pct|
+----------+----+------+------+----+------+----+-------------------+
|      3.07| 1.0|Female|   Yes| Sat|Dinner|   1|0.32573289902280134|
|      5.75| 1.0|Female|   Yes| Fri|Dinner|   2|0.17391304347826086|
|      7.25|5.15|  Male|   Yes| Sun|Dinner|   2|  0.710344827586207|
|      7.25| 1.0|Female|    No| Sat|Dinner|   1|0.13793103448275862|
|      7.51| 2.0|  Male|    No|Thur| Lunch|   2| 0.2663115845539281|
|      7.56|1.44|  Male|    No|Thur| Lunch|   2|0.19047619047619047|
|      7.74|1.44|  Male|   Yes| Sat|Dinner|   2|0.18604651162790697|
|      8.35| 1.5|Female|    No|Thur| Lunch|   2|0.17964071856287425|
|      8.51|1.25|Female|    No|Thur| Lunch|   2|0.14688601645123384|
|      8.52|1.48|  Male|    No|Thur| Lunch|   2|0.17370892018779344|
|      8.58|1.92|  Male|   Yes| Fri| Lunch|   1|0.22377622377622378|
|      8.77| 2.0|  Male|    No| Su

In [45]:
#sort by day then size
df.sort(df.day, df.size).show()

+----------+----+------+------+---+------+----+-------------------+
|total_bill| tip|   sex|smoker|day|  time|size|            tip_pct|
+----------+----+------+------+---+------+----+-------------------+
|      8.58|1.92|  Male|   Yes|Fri| Lunch|   1|0.22377622377622378|
|     16.27| 2.5|Female|   Yes|Fri| Lunch|   2|0.15365703749231716|
|     22.75|3.25|Female|    No|Fri|Dinner|   2|0.14285714285714285|
|     13.42|1.58|  Male|   Yes|Fri| Lunch|   2|0.11773472429210134|
|     16.32| 4.3|Female|   Yes|Fri|Dinner|   2|0.26348039215686275|
|     27.28| 4.0|  Male|   Yes|Fri|Dinner|   2|0.14662756598240467|
|     21.01| 3.0|  Male|   Yes|Fri|Dinner|   2| 0.1427891480247501|
|     12.46| 1.5|  Male|    No|Fri|Dinner|   2| 0.1203852327447833|
|     12.16| 2.2|  Male|   Yes|Fri| Lunch|   2|0.18092105263157895|
|     10.09| 2.0|Female|   Yes|Fri| Lunch|   2|0.19821605550049554|
|     13.42|3.48|Female|   Yes|Fri| Lunch|   2| 0.2593144560357675|
|     28.97| 3.0|  Male|   Yes|Fri|Dinner|   2|0

In [46]:
from pyspark.sql.functions import asc, desc, col

In [47]:
#sort by day, then by time (a-z), then size (largest to smallest)
df.sort(df.day, asc('time'), desc('size')).show()

+----------+----+------+------+---+------+----+-------------------+
|total_bill| tip|   sex|smoker|day|  time|size|            tip_pct|
+----------+----+------+------+---+------+----+-------------------+
|     40.17|4.73|  Male|   Yes|Fri|Dinner|   4| 0.1177495643515061|
|     22.49| 3.5|  Male|    No|Fri|Dinner|   2|0.15562472209871056|
|      5.75| 1.0|Female|   Yes|Fri|Dinner|   2|0.17391304347826086|
|     11.35| 2.5|Female|   Yes|Fri|Dinner|   2|0.22026431718061676|
|     15.38| 3.0|Female|   Yes|Fri|Dinner|   2|0.19505851755526657|
|     22.75|3.25|Female|    No|Fri|Dinner|   2|0.14285714285714285|
|     27.28| 4.0|  Male|   Yes|Fri|Dinner|   2|0.14662756598240467|
|     28.97| 3.0|  Male|   Yes|Fri|Dinner|   2|0.10355540214014498|
|     12.03| 1.5|  Male|   Yes|Fri|Dinner|   2|0.12468827930174564|
|     21.01| 3.0|  Male|   Yes|Fri|Dinner|   2| 0.1427891480247501|
|     16.32| 4.3|Female|   Yes|Fri|Dinner|   2|0.26348039215686275|
|     12.46| 1.5|  Male|    No|Fri|Dinner|   2| 

In [48]:
#use col function- nulls first
col('size').asc()

Column<'size ASC NULLS FIRST'>

In [None]:
#

In [49]:
#sort by size columns (desc:largest first) 
#then by time column (ascending: a-z)
df.sort(col('size').desc(), col('time')).show()

+----------+----+------+------+----+------+----+-------------------+
|total_bill| tip|   sex|smoker| day|  time|size|            tip_pct|
+----------+----+------+------+----+------+----+-------------------+
|     48.17| 5.0|  Male|    No| Sun|Dinner|   6|0.10379904504878555|
|      34.3| 6.7|  Male|    No|Thur| Lunch|   6|0.19533527696793004|
|     27.05| 5.0|Female|    No|Thur| Lunch|   6|0.18484288354898337|
|      29.8| 4.2|Female|    No|Thur| Lunch|   6|0.14093959731543623|
|     20.69| 5.0|  Male|    No| Sun|Dinner|   5| 0.2416626389560174|
|     29.85|5.14|Female|    No| Sun|Dinner|   5| 0.1721943048576214|
|     28.15| 3.0|  Male|   Yes| Sat|Dinner|   5|0.10657193605683837|
|     30.46| 2.0|  Male|   Yes| Sun|Dinner|   5|0.06565988181221273|
|     41.19| 5.0|  Male|    No|Thur| Lunch|   5|0.12138868657441128|
|     18.35| 2.5|  Male|    No| Sat|Dinner|   4|0.13623978201634876|
|     34.81| 5.2|Female|    No| Sun|Dinner|   4|0.14938236139040506|
|     25.29|4.71|  Male|    No| Su

<hr style="border:2px solid black"> </hr>

## Filtering

In [50]:
#df[df.tip <4] - in pandas
df.where(df.tip < 4).show() #in spark

+----------+----+------+------+---+------+----+-------------------+
|total_bill| tip|   sex|smoker|day|  time|size|            tip_pct|
+----------+----+------+------+---+------+----+-------------------+
|     16.99|1.01|Female|    No|Sun|Dinner|   2|0.05944673337257211|
|     10.34|1.66|  Male|    No|Sun|Dinner|   3|0.16054158607350097|
|     21.01| 3.5|  Male|    No|Sun|Dinner|   3|0.16658733936220846|
|     23.68|3.31|  Male|    No|Sun|Dinner|   2| 0.1397804054054054|
|     24.59|3.61|Female|    No|Sun|Dinner|   4|0.14680764538430255|
|      8.77| 2.0|  Male|    No|Sun|Dinner|   2|0.22805017103762829|
|     26.88|3.12|  Male|    No|Sun|Dinner|   4|0.11607142857142858|
|     15.04|1.96|  Male|    No|Sun|Dinner|   2|0.13031914893617022|
|     14.78|3.23|  Male|    No|Sun|Dinner|   2| 0.2185385656292287|
|     10.27|1.71|  Male|    No|Sun|Dinner|   2| 0.1665043816942551|
|     15.42|1.57|  Male|    No|Sun|Dinner|   2|0.10181582360570687|
|     18.43| 3.0|  Male|    No|Sun|Dinner|   4|0

In [51]:
#where tip is less than 4
mask = df.tip < 4
df.where(mask).show()

+----------+----+------+------+---+------+----+-------------------+
|total_bill| tip|   sex|smoker|day|  time|size|            tip_pct|
+----------+----+------+------+---+------+----+-------------------+
|     16.99|1.01|Female|    No|Sun|Dinner|   2|0.05944673337257211|
|     10.34|1.66|  Male|    No|Sun|Dinner|   3|0.16054158607350097|
|     21.01| 3.5|  Male|    No|Sun|Dinner|   3|0.16658733936220846|
|     23.68|3.31|  Male|    No|Sun|Dinner|   2| 0.1397804054054054|
|     24.59|3.61|Female|    No|Sun|Dinner|   4|0.14680764538430255|
|      8.77| 2.0|  Male|    No|Sun|Dinner|   2|0.22805017103762829|
|     26.88|3.12|  Male|    No|Sun|Dinner|   4|0.11607142857142858|
|     15.04|1.96|  Male|    No|Sun|Dinner|   2|0.13031914893617022|
|     14.78|3.23|  Male|    No|Sun|Dinner|   2| 0.2185385656292287|
|     10.27|1.71|  Male|    No|Sun|Dinner|   2| 0.1665043816942551|
|     15.42|1.57|  Male|    No|Sun|Dinner|   2|0.10181582360570687|
|     18.43| 3.0|  Male|    No|Sun|Dinner|   4|0

### OR : multiple filters at onces

In [52]:
#When meal time is dinner 
##OR 
#tip is less than/equal to 2
#sort by tip (asc: smaller to larger)

df.filter((df.time == "Dinner") | (df.tip <= 2)).sort('tip').show()

+----------+----+------+------+----+------+----+-------------------+
|total_bill| tip|   sex|smoker| day|  time|size|            tip_pct|
+----------+----+------+------+----+------+----+-------------------+
|      3.07| 1.0|Female|   Yes| Sat|Dinner|   1|0.32573289902280134|
|      12.6| 1.0|  Male|   Yes| Sat|Dinner|   2|0.07936507936507936|
|      7.25| 1.0|Female|    No| Sat|Dinner|   1|0.13793103448275862|
|      5.75| 1.0|Female|   Yes| Fri|Dinner|   2|0.17391304347826086|
|     16.99|1.01|Female|    No| Sun|Dinner|   2|0.05944673337257211|
|      12.9| 1.1|Female|   Yes| Sat|Dinner|   2|0.08527131782945736|
|     32.83|1.17|  Male|   Yes| Sat|Dinner|   2|0.03563813585135547|
|     10.51|1.25|  Male|    No| Sat|Dinner|   2|0.11893434823977164|
|     10.07|1.25|  Male|    No| Sat|Dinner|   2|0.12413108242303872|
|      8.51|1.25|Female|    No|Thur| Lunch|   2|0.14688601645123384|
|      9.68|1.32|  Male|    No| Sun|Dinner|   2|0.13636363636363638|
|     18.64|1.36|Female|    No|Thu

### and : multiple filters at once

In [53]:
#using . where 
# must be smoker, must be on saturday
df.where(df.smoker == "Yes").where(df.day == "Sat").show()

+----------+----+------+------+---+------+----+-------------------+
|total_bill| tip|   sex|smoker|day|  time|size|            tip_pct|
+----------+----+------+------+---+------+----+-------------------+
|     38.01| 3.0|  Male|   Yes|Sat|Dinner|   4|0.07892659826361484|
|     11.24|1.76|  Male|   Yes|Sat|Dinner|   2|0.15658362989323843|
|     20.29|3.21|  Male|   Yes|Sat|Dinner|   2| 0.1582060128141942|
|     13.81| 2.0|  Male|   Yes|Sat|Dinner|   2| 0.1448225923244026|
|     11.02|1.98|  Male|   Yes|Sat|Dinner|   2|0.17967332123411978|
|     18.29|3.76|  Male|   Yes|Sat|Dinner|   4|0.20557681793329688|
|      3.07| 1.0|Female|   Yes|Sat|Dinner|   1|0.32573289902280134|
|     15.01|2.09|  Male|   Yes|Sat|Dinner|   2|0.13924050632911392|
|     26.86|3.14|Female|   Yes|Sat|Dinner|   2|0.11690245718540582|
|     25.28| 5.0|Female|   Yes|Sat|Dinner|   2|0.19778481012658228|
|     17.92|3.08|  Male|   Yes|Sat|Dinner|   2|           0.171875|
|      44.3| 2.5|Female|   Yes|Sat|Dinner|   3|0

<hr style="border:2px solid black"> </hr>

## Aggregating

- .crosstab is just for counts, for other methods of summarizing groups, 

- use .groupBy (maybe in combination with .pivot) + .agg.

In [54]:
from pyspark.sql.functions import mean, min, max

In [56]:
df.show(5)

+----------+----+------+------+---+------+----+-------------------+
|total_bill| tip|   sex|smoker|day|  time|size|            tip_pct|
+----------+----+------+------+---+------+----+-------------------+
|     16.99|1.01|Female|    No|Sun|Dinner|   2|0.05944673337257211|
|     10.34|1.66|  Male|    No|Sun|Dinner|   3|0.16054158607350097|
|     21.01| 3.5|  Male|    No|Sun|Dinner|   3|0.16658733936220846|
|     23.68|3.31|  Male|    No|Sun|Dinner|   2| 0.1397804054054054|
|     24.59|3.61|Female|    No|Sun|Dinner|   4|0.14680764538430255|
+----------+----+------+------+---+------+----+-------------------+
only showing top 5 rows



In [57]:
df.groupBy('time').agg(mean('tip')).show()

+------+------------------+
|  time|          avg(tip)|
+------+------------------+
| Lunch|2.7280882352941176|
|Dinner| 3.102670454545455|
+------+------------------+



In [58]:
df.groupBy('time').agg(min('tip'), mean('tip'), max('tip')).show()

+------+--------+------------------+--------+
|  time|min(tip)|          avg(tip)|max(tip)|
+------+--------+------------------+--------+
| Lunch|    1.25|2.7280882352941176|     6.7|
|Dinner|     1.0| 3.102670454545455|    10.0|
+------+--------+------------------+--------+



In [59]:
df.groupBy('time').agg(mean('tip').alias('avg_tip')).show()

+------+------------------+
|  time|           avg_tip|
+------+------------------+
| Lunch|2.7280882352941176|
|Dinner| 3.102670454545455|
+------+------------------+



In [60]:
df.groupBy('time', 'day').agg(mean('total_bill')).show()

+------+----+------------------+
|  time| day|   avg(total_bill)|
+------+----+------------------+
| Lunch|Thur|17.664754098360653|
|Dinner|Thur|             18.78|
| Lunch| Fri|12.845714285714285|
|Dinner| Fri| 19.66333333333333|
|Dinner| Sun|21.409999999999997|
|Dinner| Sat|20.441379310344825|
+------+----+------------------+



In [61]:
df.crosstab('time', 'day').show()

+--------+---+---+---+----+
|time_day|Fri|Sat|Sun|Thur|
+--------+---+---+---+----+
|   Lunch|  7|  0|  0|  61|
|  Dinner| 12| 87| 76|   1|
+--------+---+---+---+----+



In [63]:
df.groupBy('time', 'day').agg(sum('total_bill')).sort('time','day').show()

+------+----+------------------+
|  time| day|   sum(total_bill)|
+------+----+------------------+
|Dinner| Fri|235.95999999999998|
|Dinner| Sat|1778.3999999999996|
|Dinner| Sun|1627.1599999999999|
|Dinner|Thur|             18.78|
| Lunch| Fri|             89.92|
| Lunch|Thur|           1077.55|
+------+----+------------------+



In [64]:
#pivot the above table
df.groupBy('time').pivot('day').agg(sum('total_bill')).show()

+------+------------------+------------------+------------------+-------+
|  time|               Fri|               Sat|               Sun|   Thur|
+------+------------------+------------------+------------------+-------+
| Lunch|             89.92|              null|              null|1077.55|
|Dinner|235.95999999999998|1778.3999999999996|1627.1599999999999|  18.78|
+------+------------------+------------------+------------------+-------+



In [62]:
df.groupBy('time').pivot('day').agg(mean('total_bill')).show()

+------+------------------+------------------+------------------+------------------+
|  time|               Fri|               Sat|               Sun|              Thur|
+------+------------------+------------------+------------------+------------------+
| Lunch|12.845714285714285|              null|              null|17.664754098360653|
|Dinner| 19.66333333333333|20.441379310344825|21.409999999999997|             18.78|
+------+------------------+------------------+------------------+------------------+



<hr style="border:2px solid black"> </hr>

# Additional Features

## Spark SQL
- write a sql query, get back a DF

In [65]:
#registers a sql table with the running spark
df.createOrReplaceTempView('tips')

In [66]:
spark.sql('''
SELECT *
FROM tips
''').show()

+----------+----+------+------+---+------+----+-------------------+
|total_bill| tip|   sex|smoker|day|  time|size|            tip_pct|
+----------+----+------+------+---+------+----+-------------------+
|     16.99|1.01|Female|    No|Sun|Dinner|   2|0.05944673337257211|
|     10.34|1.66|  Male|    No|Sun|Dinner|   3|0.16054158607350097|
|     21.01| 3.5|  Male|    No|Sun|Dinner|   3|0.16658733936220846|
|     23.68|3.31|  Male|    No|Sun|Dinner|   2| 0.1397804054054054|
|     24.59|3.61|Female|    No|Sun|Dinner|   4|0.14680764538430255|
|     25.29|4.71|  Male|    No|Sun|Dinner|   4|0.18623962040332148|
|      8.77| 2.0|  Male|    No|Sun|Dinner|   2|0.22805017103762829|
|     26.88|3.12|  Male|    No|Sun|Dinner|   4|0.11607142857142858|
|     15.04|1.96|  Male|    No|Sun|Dinner|   2|0.13031914893617022|
|     14.78|3.23|  Male|    No|Sun|Dinner|   2| 0.2185385656292287|
|     10.27|1.71|  Male|    No|Sun|Dinner|   2| 0.1665043816942551|
|     35.26| 5.0|Female|    No|Sun|Dinner|   4|0

In [69]:
# find the tip, total_bill, and day with the highest overall sales for that day
spark.sql('''
SELECT tip, total_bill, day
FROM tips
WHERE day = (
    SELECT day
    FROM tips
    GROUP BY day
    ORDER BY sum(total_bill) DESC
    LIMIT 1
)    
''').show()

+----+----------+---+
| tip|total_bill|day|
+----+----------+---+
|3.35|     20.65|Sat|
|4.08|     17.92|Sat|
|2.75|     20.29|Sat|
|2.23|     15.77|Sat|
|7.58|     39.42|Sat|
|3.18|     19.82|Sat|
|2.34|     17.81|Sat|
| 2.0|     13.37|Sat|
| 2.0|     12.69|Sat|
| 4.3|      21.7|Sat|
| 3.0|     19.65|Sat|
|1.45|      9.55|Sat|
| 2.5|     18.35|Sat|
| 3.0|     15.06|Sat|
|2.45|     20.69|Sat|
|3.27|     17.78|Sat|
| 3.6|     24.06|Sat|
| 2.0|     16.31|Sat|
|3.07|     16.93|Sat|
|2.31|     18.69|Sat|
+----+----------+---+
only showing top 20 rows



<hr style="border:2px solid black"> </hr>

# More Spark Dataframe Manipulation

## Mixing in SQL expressions
- expr()

In [72]:
from pyspark.sql.functions import expr

In [73]:
df.select(
    '*',
    expr('tip / total_bill as tip_pct')
).where(
    expr('day = "Sun" AND time = "Dinner"')
).show()

+----------+----+------+------+---+------+----+-------------------+-------------------+
|total_bill| tip|   sex|smoker|day|  time|size|            tip_pct|            tip_pct|
+----------+----+------+------+---+------+----+-------------------+-------------------+
|     16.99|1.01|Female|    No|Sun|Dinner|   2|0.05944673337257211|0.05944673337257211|
|     10.34|1.66|  Male|    No|Sun|Dinner|   3|0.16054158607350097|0.16054158607350097|
|     21.01| 3.5|  Male|    No|Sun|Dinner|   3|0.16658733936220846|0.16658733936220846|
|     23.68|3.31|  Male|    No|Sun|Dinner|   2| 0.1397804054054054| 0.1397804054054054|
|     24.59|3.61|Female|    No|Sun|Dinner|   4|0.14680764538430255|0.14680764538430255|
|     25.29|4.71|  Male|    No|Sun|Dinner|   4|0.18623962040332148|0.18623962040332148|
|      8.77| 2.0|  Male|    No|Sun|Dinner|   2|0.22805017103762829|0.22805017103762829|
|     26.88|3.12|  Male|    No|Sun|Dinner|   4|0.11607142857142858|0.11607142857142858|
|     15.04|1.96|  Male|    No|S

<hr style="border:2px solid black"> </hr>

## Joins

In [74]:
df1 = spark.createDataFrame(pd.DataFrame({
    'id': np.arange(100) + 1,
    'x': np.random.randn(100).round(3),
    'group_id': np.random.choice(range(1, 7), 100),
}))
df2 = spark.createDataFrame(pd.DataFrame({
    'id': range(1, 7),
    'group': list('abcdef')
}))
df1.show(5)
df2.show()

+---+------+--------+
| id|     x|group_id|
+---+------+--------+
|  1| 0.527|       5|
|  2| 1.287|       1|
|  3|-1.209|       3|
|  4| 0.688|       4|
|  5|-0.157|       6|
+---+------+--------+
only showing top 5 rows

+---+-----+
| id|group|
+---+-----+
|  1|    a|
|  2|    b|
|  3|    c|
|  4|    d|
|  5|    e|
|  6|    f|
+---+-----+



In [77]:
#this works but its wierd, need a condition
df1.join(df2).show()

+---+------+--------+---+-----+
| id|     x|group_id| id|group|
+---+------+--------+---+-----+
|  1| 0.527|       5|  1|    a|
|  2| 1.287|       1|  1|    a|
|  3|-1.209|       3|  1|    a|
|  4| 0.688|       4|  1|    a|
|  5|-0.157|       6|  1|    a|
|  6| 1.981|       3|  1|    a|
|  7|-2.311|       4|  1|    a|
|  8| 0.583|       2|  1|    a|
|  9|-0.751|       2|  1|    a|
| 10|-0.502|       6|  1|    a|
| 11| 0.012|       5|  1|    a|
| 12|-0.229|       3|  1|    a|
|  1| 0.527|       5|  2|    b|
|  2| 1.287|       1|  2|    b|
|  3|-1.209|       3|  2|    b|
|  4| 0.688|       4|  2|    b|
|  5|-0.157|       6|  2|    b|
|  6| 1.981|       3|  2|    b|
|  7|-2.311|       4|  2|    b|
|  8| 0.583|       2|  2|    b|
+---+------+--------+---+-----+
only showing top 20 rows



In [75]:
df_merged = df1.join(df2, df1.group_id == df2.id)
df_merged.show(5)

+---+------+--------+---+-----+
| id|     x|group_id| id|group|
+---+------+--------+---+-----+
|  5|-0.157|       6|  6|    f|
| 10|-0.502|       6|  6|    f|
| 15| 0.396|       6|  6|    f|
| 23| 1.379|       6|  6|    f|
| 27| 1.776|       6|  6|    f|
+---+------+--------+---+-----+
only showing top 5 rows



In [76]:
#rename the columns that are the same to make more legible
#now join on group_id
df1.join(df2.withColumnRenamed('id', 'group_id'), 'group_id').show(5)

+--------+---+------+-----+
|group_id| id|     x|group|
+--------+---+------+-----+
|       6|  5|-0.157|    f|
|       6| 10|-0.502|    f|
|       6| 15| 0.396|    f|
|       6| 23| 1.379|    f|
|       6| 27| 1.776|    f|
+--------+---+------+-----+
only showing top 5 rows



<hr style="border:2px solid black"> </hr>

## .explain

In [80]:
#shows how df came into existance
#read from bottom to top
df.explain()

== Physical Plan ==
*(1) Project [total_bill#0, tip#1, sex#2, smoker#3, day#4, time#5, size#6L, (tip#1 / total_bill#0) AS tip_pct#432]
+- *(1) Scan ExistingRDD[total_bill#0,tip#1,sex#2,smoker#3,day#4,time#5,size#6L]




In [81]:
#create the dataframe
df = spark.createDataFrame(pydataset.data('tips'))

In [82]:
#grab the existing dataframe and filter time for only lunch
df.filter(df.time == 'Lunch').explain()

== Physical Plan ==
*(1) Filter (isnotnull(time#1906) AND (time#1906 = Lunch))
+- *(1) Scan ExistingRDD[total_bill#1901,tip#1902,sex#1903,smoker#1904,day#1905,time#1906,size#1907L]




In [83]:
#filter lunch and size greater than 1
df.filter(df.time == 'Lunch').filter(df.size >1).explain()

== Physical Plan ==
*(1) Filter (((isnotnull(time#1906) AND isnotnull(size#1907L)) AND (time#1906 = Lunch)) AND (size#1907L > 1))
+- *(1) Scan ExistingRDD[total_bill#1901,tip#1902,sex#1903,smoker#1904,day#1905,time#1906,size#1907L]




In [84]:
#lunch, larger than 1, tip percent column
df.filter(df.time == 'Lunch').filter(df.size >1).select('*', expr('tip/ total_bill as tip_pct')).explain()

== Physical Plan ==
*(1) Project [total_bill#1901, tip#1902, sex#1903, smoker#1904, day#1905, time#1906, size#1907L, (tip#1902 / total_bill#1901) AS tip_pct#1915]
+- *(1) Filter (((isnotnull(time#1906) AND isnotnull(size#1907L)) AND (time#1906 = Lunch)) AND (size#1907L > 1))
   +- *(1) Scan ExistingRDD[total_bill#1901,tip#1902,sex#1903,smoker#1904,day#1905,time#1906,size#1907L]




In [86]:
# Scan, Filter, Project
df.filter(df.smoker =='Yes').select('smoker', 'total_bill', expr('tip/ total_bill as tip_pct')).explain()

== Physical Plan ==
*(1) Project [smoker#1904, total_bill#1901, (tip#1902 / total_bill#1901) AS tip_pct#1925]
+- *(1) Filter (isnotnull(smoker#1904) AND (smoker#1904 = Yes))
   +- *(1) Scan ExistingRDD[total_bill#1901,tip#1902,sex#1903,smoker#1904,day#1905,time#1906,size#1907L]




In [87]:
#Scan, Filer, Project
#even though order was flipped (from above cell), Spark is smart enough to fix order
#Filter before Project
df.select('smoker', 'total_bill', expr('tip/ total_bill as tip_pct')).filter(df.smoker =='Yes').explain()

== Physical Plan ==
*(1) Project [smoker#1904, total_bill#1901, (tip#1902 / total_bill#1901) AS tip_pct#1929]
+- *(1) Filter (isnotnull(smoker#1904) AND (smoker#1904 = Yes))
   +- *(1) Scan ExistingRDD[total_bill#1901,tip#1902,sex#1903,smoker#1904,day#1905,time#1906,size#1907L]




In [70]:
df.where(
    df.time == 'Dinner'
).select(
    '*',
    (df.tip / df.total_bill).alias('tip_pct'),
).explain()

== Physical Plan ==
*(1) Project [total_bill#0, tip#1, sex#2, smoker#3, day#4, time#5, size#6L, (tip#1 / total_bill#0) AS tip_pct#432, (tip#1 / total_bill#0) AS tip_pct#1680]
+- *(1) Filter (isnotnull(time#5) AND (time#5 = Dinner))
   +- *(1) Scan ExistingRDD[total_bill#0,tip#1,sex#2,smoker#3,day#4,time#5,size#6L]




In [71]:
df.select(
    '*',
    (df.tip / df.total_bill).alias('tip_pct'),
).where(
    df.time == 'Dinner'
).explain()

== Physical Plan ==
*(1) Project [total_bill#0, tip#1, sex#2, smoker#3, day#4, time#5, size#6L, (tip#1 / total_bill#0) AS tip_pct#432, (tip#1 / total_bill#0) AS tip_pct#1690]
+- *(1) Filter (isnotnull(time#5) AND (time#5 = Dinner))
   +- *(1) Scan ExistingRDD[total_bill#0,tip#1,sex#2,smoker#3,day#4,time#5,size#6L]


