{"metadata":{"kernelspec":{"language":"python","display_name":"Python 3","name":"python3"},"language_info":{"pygments_lexer":"ipython3","nbconvert_exporter":"python","version":"3.6.4","file_extension":".py","codemirror_mode":{"name":"ipython","version":3},"name":"python","mimetype":"text/x-python"}},"nbformat_minor":4,"nbformat":4,"cells":[{"cell_type":"markdown","source":"<a id=\"0\"></a>\n# <p style=\"background-color:#281F2F;height: 60px;text-align: center;vertical-align: middle;line-height: 60px;;font-family:helvetica;color:#FFFFFF;font-size:120%;text-align:center;border-radius:12px 12px;\">Libraries </p>","metadata":{"_kg_hide-output":true,"_kg_hide-input":true}},{"cell_type":"code","source":"import numpy as np \nimport pandas as pd \nimport os","metadata":{"_uuid":"8f2839f25d086af736a60e9eeb907d3b93b6e0e5","_cell_guid":"b1076dfc-b9ad-4769-8c92-a6c4dae69d19","_kg_hide-input":true,"execution":{"iopub.status.busy":"2022-06-02T18:06:35.651646Z","iopub.execute_input":"2022-06-02T18:06:35.652888Z","iopub.status.idle":"2022-06-02T18:06:35.68229Z","shell.execute_reply.started":"2022-06-02T18:06:35.652381Z","shell.execute_reply":"2022-06-02T18:06:35.681242Z"},"trusted":true},"execution_count":null,"outputs":[]},{"cell_type":"markdown","source":"<!-- <img src=\"https://raw.githubusercontent.com/paulojunqueira/files/main/101%20pyspark.jpg\">  -->\n <img src=\"https://raw.githubusercontent.com/paulojunqueira/files/main/101.gif\" style = 'class=center; margin-left: auto; margin-right: auto;width: 50%;display: block;'>  \n\n<p style = \"font-size:140%\"> This is an introductory Pyspark hands on tutorial. I have decided to study and summarize some basic concepts of pyspark functions. Hope these can help some else!\n</p>","metadata":{}},{"cell_type":"markdown","source":"\n\n\n# CONTENT TABLE\n----------------\n## Concepts\n* [1) What is Spark \\ PySpark](#01)\n* [2) Characteristics](#03)\n  \n\n## Practice\n* [1) GETTING RUNNING ON NOTEBOOKS ( Jupyter\\Colab\\Kaggle)](#1)\n* [2) READING DATA FROM DIFFERENT SOURCES](#2)\n* [3) DISPLAYING](#3)\n* [4) FILTERING](#4)\n* [5) FILTERING PART 2](#5)\n* [6) MANIPULATING DATA](#6)\n* [7) GROUPING BY](#7)\n* [8) WINDOW FUNCTION](#8)\n* [9) JOIN](#9)\n* [10) UDF](#10)\n* [11) REFERENCES](#20)\n\n## Next\n* 102 Pyspark - MLlib (Soon!) ","metadata":{}},{"cell_type":"markdown","source":"<a id=\"01\"></a>\n# <p style=\"background-color:#281F2F;height: 60px;text-align: center;vertical-align: middle;line-height: 60px;;font-family:helvetica;color:#FFFFFF;font-size:120%;text-align:center;border-radius:12px 12px;\"> What is Spark \\ PySpark </p>","metadata":{}},{"cell_type":"markdown","source":"<p style = \"font-size:120%\"> \n    Apache Spark is an open source unified computing engine for distributed data processing on computer clusters, allowing an easy and scalable developments. It supports many programing languages such as Python, R and Scala. Also, includes libraries from SQL and Machine Learning. The structure is designed to deal with huge amount of data distributing large datasets among \"computers\". Each one of the computers os called a \"Executors\" and the managemnt of all data is done by a so called \"Driver\".<br>\n    <br>\n    Pyspark is a library for python that enables to run python applications in the Apache Spark architeture, hence, allowing parallel distribuiton to cope with large data \n\n</p>\n    \n<img src=\"https://raw.githubusercontent.com/paulojunqueira/files/main/Sem%20t%C3%ADtulo.gif\" style = 'class=center; margin-left: auto; margin-right: auto;width: 50%;display: block;'>  \n\n\n ","metadata":{}},{"cell_type":"markdown","source":"<a id=\"03\"></a>\n# <p style=\"background-color:#281F2F;height: 60px;text-align: center;vertical-align: middle;line-height: 60px;;font-family:helvetica;color:#FFFFFF;font-size:120%;text-align:center;border-radius:12px 12px;\"> Characteristics </p>\n\n<p style = \"font-size:120%\">\nAmong the characteristics of Spark, some interesting ones are:\n    \n<li style = \"font-size:110%\">Fast Processing: High data processing speed that could reach 100x faster in memory and 10x faster on the disk\n<li style = \"font-size:110%\">Lazy Evaluation: The transformations applied are not done right way. Spark order all transformations in an efficient way under the hood.\n<li style = \"font-size:110%\">Support Multiple High level languages: Python, R, Scala, Java, SQL.\n<li style = \"font-size:110%\">Native Libraries: Spark has Machine Learning and Graph libraries that cope with distributed systems (MLlib, SQL, Dataframes, GraphFrames)\n\n</p>","metadata":{}},{"cell_type":"markdown","source":"--------------------","metadata":{}},{"cell_type":"markdown","source":"<a id=\"1\"></a>\n# <p style=\"background-color:#281F2F;height: 60px;text-align: center;vertical-align: middle;line-height: 60px;;font-family:helvetica;color:#FFFFFF;font-size:120%;text-align:center;border-radius:12px 12px;\">GETTING RUNNING ON NOTEBOOKS ( Jupyter\\Colab\\Kaggle)</p>","metadata":{}},{"cell_type":"code","source":"!pip install pyspark -q","metadata":{"execution":{"iopub.status.busy":"2022-06-02T18:06:35.72114Z","iopub.execute_input":"2022-06-02T18:06:35.721956Z","iopub.status.idle":"2022-06-02T18:07:31.662266Z","shell.execute_reply.started":"2022-06-02T18:06:35.721916Z","shell.execute_reply":"2022-06-02T18:07:31.660847Z"},"trusted":true},"execution_count":null,"outputs":[]},{"cell_type":"code","source":"# To run pypsark in notebooks: \nfrom pyspark import SparkContext, SparkConf\nfrom pyspark.sql import SparkSession\n\n#Import libraries to create a dataframe\n#It is also common to import as T: #from pypsark.sql.types as T\nfrom pyspark.sql.types import StructType,StructField, StringType, IntegerType\n\n\nsc = SparkContext.getOrCreate(SparkConf().setMaster(\"local[*]\"))\nspark = SparkSession.builder.getOrCreate()","metadata":{"execution":{"iopub.status.busy":"2022-06-02T18:07:31.665707Z","iopub.execute_input":"2022-06-02T18:07:31.66602Z","iopub.status.idle":"2022-06-02T18:07:37.585764Z","shell.execute_reply.started":"2022-06-02T18:07:31.665981Z","shell.execute_reply":"2022-06-02T18:07:37.584073Z"},"trusted":true},"execution_count":null,"outputs":[]},{"cell_type":"markdown","source":"<a id=\"2\"></a>\n# <p style=\"background-color:#281F2F;height: 60px;text-align: center;vertical-align: middle;line-height: 60px;;font-family:helvetica;color:#FFFFFF;font-size:120%;text-align:center;border-radius:12px 12px;\"> READING DATA FROM DIFFERENT SOURCES</p>\n\n<li style = \"font-size:110%\">Create a dataframe in Pyspark\n<li style = \"font-size:110%\">Read a csv in Pyspark\n<li style = \"font-size:110%\">Convert a Pandas df to Pyspark df","metadata":{}},{"cell_type":"code","source":"\n#Toy dataset for example\nexample_data = [(100,\"Brazil\",\"1000\",\"A100\",6),\n    (101,\"Spain\",\"2000\",\"MA100\",2),\n    (102,\"EUA\",\"3000\",\"A200\",10),\n    (110,\"Mexico\",\"B100\",\"F400\",8),\n    (200,\"Japan\",\"5000\",\"A100\",9),\n    (880,\"EUA\",\"500\",\"Z120\",1)\n  ]\n\n#Defining a Schema of the Dataset (name, type of column and if its nullable)\nschema = StructType([StructField('key', IntegerType(), True),\n                    StructField('C1', StringType(), True),\n                    StructField('C2', StringType(), True),\n                    StructField('C3', StringType(), True),\n                    StructField('C4', IntegerType(), True)])\n\ndf_example = spark.createDataFrame(example_data, schema = schema)\n\n#displaying data \ndf_example.show()","metadata":{"execution":{"iopub.status.busy":"2022-06-02T18:55:45.723975Z","iopub.execute_input":"2022-06-02T18:55:45.724258Z","iopub.status.idle":"2022-06-02T18:55:45.883904Z","shell.execute_reply.started":"2022-06-02T18:55:45.724229Z","shell.execute_reply":"2022-06-02T18:55:45.883002Z"},"trusted":true},"execution_count":null,"outputs":[]},{"cell_type":"code","source":"#Toy dataset for example\nexample_data_2 = [(100,\"BB\",10),\n    (101,\"XX\",99),\n    (102,\"AN\",898),\n    (110,\"AC\",567),\n    (200,\"AV\",344),\n    (300,\"FV\",834),\n    (111,\"ZW\",54)\n  ]\n\n\n#Defining a Schema of the Dataset (name, type of column and if its nullable)\nschema_2 = StructType([StructField('key', IntegerType(), True),\n                    StructField('C5', StringType(), True),\n                    StructField('C6', IntegerType(), True)])\n\ndf_example_2 = spark.createDataFrame(example_data_2, schema = schema_2)\n\n#displaying data \ndf_example_2.show()","metadata":{"execution":{"iopub.status.busy":"2022-06-02T18:55:46.664975Z","iopub.execute_input":"2022-06-02T18:55:46.665259Z","iopub.status.idle":"2022-06-02T18:55:46.764046Z","shell.execute_reply.started":"2022-06-02T18:55:46.66523Z","shell.execute_reply":"2022-06-02T18:55:46.763152Z"},"trusted":true},"execution_count":null,"outputs":[]},{"cell_type":"code","source":"# 2) reading a csv in spark\ndf_t = spark.read.csv('../input/titanic/train.csv', header = True, inferSchema = True)\n\n#displaying data\ndf_t.show(5)","metadata":{"execution":{"iopub.status.busy":"2022-06-02T18:07:44.646898Z","iopub.execute_input":"2022-06-02T18:07:44.647229Z","iopub.status.idle":"2022-06-02T18:07:46.785903Z","shell.execute_reply.started":"2022-06-02T18:07:44.647196Z","shell.execute_reply":"2022-06-02T18:07:46.784918Z"},"trusted":true},"execution_count":null,"outputs":[]},{"cell_type":"code","source":"# 3) Read a pandas dataframe \ndf_pandas = pd.read_csv('../input/wine-quality-dataset/WineQT.csv')\n\n#Converting it to a a spark dataframe\ndf = spark.createDataFrame(df_pandas)\n\n#Going back to pandas (warning, for huge datasets in distributed systems this command may cause memory overload)\n# df.toPandas()\n\n#displaying data\ndf.show(5)","metadata":{"execution":{"iopub.status.busy":"2022-06-02T18:07:46.78723Z","iopub.execute_input":"2022-06-02T18:07:46.787586Z","iopub.status.idle":"2022-06-02T18:07:47.346415Z","shell.execute_reply.started":"2022-06-02T18:07:46.78754Z","shell.execute_reply":"2022-06-02T18:07:47.345401Z"},"trusted":true},"execution_count":null,"outputs":[]},{"cell_type":"markdown","source":"<a id=\"3\"></a>\n# <p style=\"background-color:#281F2F;height: 60px;text-align: center;vertical-align: middle;line-height: 60px;;font-family:helvetica;color:#FFFFFF;font-size:120%;text-align:center;border-radius:12px 12px;\"> DISPLAYING</p>\n\n<li style = \"font-size:110%\"> Displaying table\n<li style = \"font-size:110%\"> Analysing statistics and types\n<li style = \"font-size:110%\"> Selecting variables\n<li style = \"font-size:110%\"> Droping Variables","metadata":{}},{"cell_type":"code","source":"# As Already seen, showing the results. Also could use a parameter to change the number of rows displayed\nprint('Example_1')\ndf.show(10)\n\nprint('Example_2')\ndf.show(5)","metadata":{"execution":{"iopub.status.busy":"2022-06-02T18:07:47.347936Z","iopub.execute_input":"2022-06-02T18:07:47.34826Z","iopub.status.idle":"2022-06-02T18:07:47.700442Z","shell.execute_reply.started":"2022-06-02T18:07:47.348214Z","shell.execute_reply":"2022-06-02T18:07:47.69938Z"},"trusted":true},"execution_count":null,"outputs":[]},{"cell_type":"code","source":"# Use limit() to determine the number of lines to show as well\ndf.limit(2).show()","metadata":{"execution":{"iopub.status.busy":"2022-06-02T18:07:47.703676Z","iopub.execute_input":"2022-06-02T18:07:47.705674Z","iopub.status.idle":"2022-06-02T18:07:48.635498Z","shell.execute_reply.started":"2022-06-02T18:07:47.705615Z","shell.execute_reply":"2022-06-02T18:07:48.634562Z"},"trusted":true},"execution_count":null,"outputs":[]},{"cell_type":"code","source":"#Checking schema of df ( types of variables )\ndf.printSchema()","metadata":{"execution":{"iopub.status.busy":"2022-06-02T18:07:48.636963Z","iopub.execute_input":"2022-06-02T18:07:48.63728Z","iopub.status.idle":"2022-06-02T18:07:48.652495Z","shell.execute_reply.started":"2022-06-02T18:07:48.637243Z","shell.execute_reply":"2022-06-02T18:07:48.651649Z"},"trusted":true},"execution_count":null,"outputs":[]},{"cell_type":"code","source":"# Checking Columns' name\ndf.columns","metadata":{"execution":{"iopub.status.busy":"2022-06-02T18:07:48.653752Z","iopub.execute_input":"2022-06-02T18:07:48.654863Z","iopub.status.idle":"2022-06-02T18:07:48.672924Z","shell.execute_reply.started":"2022-06-02T18:07:48.654816Z","shell.execute_reply":"2022-06-02T18:07:48.671804Z"},"trusted":true},"execution_count":null,"outputs":[]},{"cell_type":"code","source":"# Checking statistics of these variables\ndf.select('fixed acidity', 'quality').describe().show()","metadata":{"execution":{"iopub.status.busy":"2022-06-02T18:07:48.680304Z","iopub.execute_input":"2022-06-02T18:07:48.68074Z","iopub.status.idle":"2022-06-02T18:07:49.574702Z","shell.execute_reply.started":"2022-06-02T18:07:48.680671Z","shell.execute_reply":"2022-06-02T18:07:49.573548Z"},"trusted":true},"execution_count":null,"outputs":[]},{"cell_type":"code","source":"# Counting the number of rows (samples)\nprint('Number of rows: ',df.count())\n\n#Counting the number of columns (variables)\nprint('Number of Columns: ', len(df.columns))","metadata":{"execution":{"iopub.status.busy":"2022-06-02T18:07:49.576065Z","iopub.execute_input":"2022-06-02T18:07:49.576753Z","iopub.status.idle":"2022-06-02T18:07:49.870814Z","shell.execute_reply.started":"2022-06-02T18:07:49.576668Z","shell.execute_reply":"2022-06-02T18:07:49.869897Z"},"trusted":true},"execution_count":null,"outputs":[]},{"cell_type":"code","source":"#Cross table between variables\ndf.crosstab('pH', 'quality').show()","metadata":{"execution":{"iopub.status.busy":"2022-06-02T18:07:49.872143Z","iopub.execute_input":"2022-06-02T18:07:49.872798Z","iopub.status.idle":"2022-06-02T18:07:51.104296Z","shell.execute_reply.started":"2022-06-02T18:07:49.872749Z","shell.execute_reply":"2022-06-02T18:07:51.103375Z"},"trusted":true},"execution_count":null,"outputs":[]},{"cell_type":"code","source":"# Selecting specific columns\ndf.select('fixed acidity', 'quality').show(5)","metadata":{"execution":{"iopub.status.busy":"2022-06-02T18:07:51.105432Z","iopub.execute_input":"2022-06-02T18:07:51.105707Z","iopub.status.idle":"2022-06-02T18:07:51.320537Z","shell.execute_reply.started":"2022-06-02T18:07:51.105672Z","shell.execute_reply":"2022-06-02T18:07:51.319776Z"},"trusted":true},"execution_count":null,"outputs":[]},{"cell_type":"code","source":"#Ordering by a column\ndf.orderBy('quality', ascending = False).show(10)\n\n#Ordering by a column\ndf.orderBy('quality', ascending = True).show(10)","metadata":{"execution":{"iopub.status.busy":"2022-06-02T18:07:51.321922Z","iopub.execute_input":"2022-06-02T18:07:51.322657Z","iopub.status.idle":"2022-06-02T18:07:51.955178Z","shell.execute_reply.started":"2022-06-02T18:07:51.322616Z","shell.execute_reply":"2022-06-02T18:07:51.954264Z"},"trusted":true},"execution_count":null,"outputs":[]},{"cell_type":"code","source":"#Removing variables (columns)\ndf.drop('fixed acidity','volatile acidity').show(5)","metadata":{"execution":{"iopub.status.busy":"2022-06-02T18:07:51.956359Z","iopub.execute_input":"2022-06-02T18:07:51.956656Z","iopub.status.idle":"2022-06-02T18:07:52.180761Z","shell.execute_reply.started":"2022-06-02T18:07:51.956617Z","shell.execute_reply":"2022-06-02T18:07:52.179625Z"},"trusted":true},"execution_count":null,"outputs":[]},{"cell_type":"markdown","source":"<a id=\"4\"></a>\n# <p style=\"background-color:#281F2F;height: 60px;text-align: center;vertical-align: middle;line-height: 60px;;font-family:helvetica;color:#FFFFFF;font-size:120%;text-align:center;border-radius:12px 12px;\"> FILTERING</p>\n<li style = \"font-size:110%\">To filter data by some conditions, we are going to use the SQL funtions.  \n<li style = \"font-size:110%\">The sql.functions is going to call many sql functions that are very useful when dealing with data.To check the whole possibility:\n    \n[SQL Functions in Spark](https://spark.apache.org/docs/2.4.0/api/python/pyspark.sql.html#pyspark.sql.functions.arrays_zip)\n","metadata":{}},{"cell_type":"code","source":"import pyspark.sql.functions as F\n\n#filtering Data\nprint('---- Example_1 ----')\ndf.filter(F.col('quality')>5).show(5)\n\n#filtering Data\nprint('---- Example_2 ----')\ndf.filter(F.col('fixed acidity')< 5).show(5)\n\n#filtering Data\nprint('---- Example_3 Multiple conditions ----')\ndf.filter((F.col('fixed acidity') < 10) & (F.col('quality') > 5)).show(5)\n\n#filtering Data\nprint('---- Example_4 Multiple conditions ----')\ndf.filter((F.col('quality') > 4) | (F.col('quality') < 8)).show(5)\n\n#filtering Data\nprint('Example_5 in list')\ndf.filter((F.col('quality').isin([5,6,7]))).show(5)\n","metadata":{"execution":{"iopub.status.busy":"2022-06-02T18:07:52.182182Z","iopub.execute_input":"2022-06-02T18:07:52.182538Z","iopub.status.idle":"2022-06-02T18:07:53.258969Z","shell.execute_reply.started":"2022-06-02T18:07:52.182489Z","shell.execute_reply":"2022-06-02T18:07:53.257809Z"},"trusted":true},"execution_count":null,"outputs":[]},{"cell_type":"markdown","source":"<a id=\"5\"></a>\n# <p style=\"background-color:#281F2F;height: 60px;text-align: center;vertical-align: middle;line-height: 60px;;font-family:helvetica;color:#FFFFFF;font-size:120%;text-align:center;border-radius:12px 12px;\"> FILTERING PART 2</p>","metadata":{}},{"cell_type":"code","source":"# Part 2\n\n#filtering Data\nprint('---- Example_6 Categorical ----')\ndf_t.filter((F.col('Embarked') ==  'S') & (F.col('Sex') ==  'female')).show(5)\n\n#filtering Data\nprint('---- Example_7 Exclusion ----')\ndf_t.filter(~(F.col('Embarked') ==  'S')).show(5)\n\n#filtering Data\nprint('---- Example_8 Exclusion SQL  ----')\ndf_t.filter(('Embarked !=  \"C\"')).show(5)\ndf_t.filter(('Embarked <>  \"Q\"')).show(5)\n\n#filtering Data\nprint('---- Example_9 Cabin Starts with A----')\ndf_t.filter((F.col('Cabin').startswith('A'))).show(5)\n\n#filtering Data\nprint('---- Example_10 Cabin Ends with 3 ----')\ndf_t.filter((F.col('Cabin').endswith('3'))).show(5)\n\n#filtering Data\nprint('---- Example_11 Name Contains \"Miss\" ----')\ndf_t.filter((F.col('Name').contains('Miss'))).show(5)\n\n#filtering Data\nprint('---- Example_12 like (REGEX)----')\ndf_t.filter((F.col('Name').like('%John%'))).show(5)\n\n#filtering Data\nprint('---- Example_13 Cabin null----')\ndf_t.filter((F.col('Cabin').isNull())).show(5)\n\n#filtering Data\nprint('---- Example_14 Cabin not null----')\ndf_t.filter((F.col('Cabin').isNotNull())).show(5)","metadata":{"execution":{"iopub.status.busy":"2022-06-02T18:07:53.260554Z","iopub.execute_input":"2022-06-02T18:07:53.260928Z","iopub.status.idle":"2022-06-02T18:07:55.290635Z","shell.execute_reply.started":"2022-06-02T18:07:53.260881Z","shell.execute_reply":"2022-06-02T18:07:55.289481Z"},"trusted":true},"execution_count":null,"outputs":[]},{"cell_type":"markdown","source":"<a id=\"6\"></a>\n# <p style=\"background-color:#281F2F;height: 60px;text-align: center;vertical-align: middle;line-height: 60px;;font-family:helvetica;color:#FFFFFF;font-size:120%;text-align:center;border-radius:12px 12px;\"> MANIPULATING DATA</p>\n<li style = \"font-size:110%\"> Create new columns \n<li style = \"font-size:110%\"> Changing columns \n\nwithColumn() F.round()  F.when() ","metadata":{}},{"cell_type":"code","source":"#Manipulating Data\nprint('---- Example_1 Adding up two Columns----') #Same goes to all math operators\ndf.withColumn('VOLATILE + CITRIC', F.col('volatile acidity') + F.col('citric acid')).show(5)\n\n#Manipulating Data\nprint('---- Example_2 Create a column of ones----') \ndf.withColumn('DUMMY_ONES', F.lit(1)).show(5)\n\n#Manipulating Data\nprint('---- Example_3 round pH variable----') \ndf.withColumn('pH_ROUNDED', F.round('pH',1)).show(5)","metadata":{"execution":{"iopub.status.busy":"2022-06-02T18:07:55.292123Z","iopub.execute_input":"2022-06-02T18:07:55.294215Z","iopub.status.idle":"2022-06-02T18:07:55.771827Z","shell.execute_reply.started":"2022-06-02T18:07:55.294159Z","shell.execute_reply":"2022-06-02T18:07:55.77095Z"},"trusted":true},"execution_count":null,"outputs":[]},{"cell_type":"markdown","source":"- Conditionals (F.when('Condition', value ))","metadata":{}},{"cell_type":"code","source":"#Manipulating Data\n\nprint('---- Example_1 Conditional based on Survived Column----') \ndf_t.withColumn('Survived_TEXT', F.when(F.col('Survived') == 0, 'NO')\\\n                                 .otherwise('YES')).show(5)\n\nprint('---- Example_2 Multiple Conditions----') \ndf_t.withColumn('COD_EMBARKED', F.when(F.col('Embarked') == 'C', 'YES_C')\\\n                                 .when(F.col('Embarked') == 'S', 'NO_S')\\\n                                 .otherwise('MAYBE_OTHER')).show(5)","metadata":{"execution":{"iopub.status.busy":"2022-06-02T18:07:55.773071Z","iopub.execute_input":"2022-06-02T18:07:55.773393Z","iopub.status.idle":"2022-06-02T18:07:56.204353Z","shell.execute_reply.started":"2022-06-02T18:07:55.773351Z","shell.execute_reply":"2022-06-02T18:07:56.20348Z"},"trusted":true},"execution_count":null,"outputs":[]},{"cell_type":"markdown","source":"<li style = \"font-size:110%\">Changing type of variables (importing the sql types module as T)\n<li style = \"font-size:110%\">DoubleType, IntegerType, StringType...","metadata":{}},{"cell_type":"code","source":"import  pyspark.sql.types as T \n\n#Manipulating Data\nprint('Original Schema')\ndf_t.select('Survived','Age').printSchema()\n\n#Manipulating Data\nprint('---- Example_1 Survid to type Double----') \ndf_t_1 = df_t.withColumn('Survived', F.col('Survived').cast(T.DoubleType()))\ndf_t_1.show(5)\nprint('New Survived Schema')\ndf_t_1.select('Survived').printSchema()\n\n#Manipulating Data\nprint('---- Example_2 Age to type String ----') \ndf_t_1 = df_t.withColumn('Age', F.col('Age').cast(T.StringType()))\ndf_t_1.show(5)\nprint('New Age Schema')\ndf_t_1.select('Age').printSchema()","metadata":{"execution":{"iopub.status.busy":"2022-06-02T18:07:56.205495Z","iopub.execute_input":"2022-06-02T18:07:56.205789Z","iopub.status.idle":"2022-06-02T18:07:56.859091Z","shell.execute_reply.started":"2022-06-02T18:07:56.205745Z","shell.execute_reply":"2022-06-02T18:07:56.858084Z"},"trusted":true},"execution_count":null,"outputs":[]},{"cell_type":"markdown","source":"<a id=\"7\"></a>\n# <p style=\"background-color:#281F2F;height: 60px;text-align: center;vertical-align: middle;line-height: 60px;;font-family:helvetica;color:#FFFFFF;font-size:120%;text-align:center;border-radius:12px 12px;\"> GROUPING BY</p>","metadata":{}},{"cell_type":"code","source":"#Group By\nprint('---- Example_1 ---- Grouping by Sex and counting the total samples of each Genre')\ndf_t.groupBy('Sex').count().show()\n\n#Group By\nprint('---- Example_2 ---- Grouping by Survived and counting the total samples of each class')\ndf_t.groupBy('Survived').count().show()\n\n#Group By\nprint('---- Example_3 ---- Grouping by Sex and survived and counting the total samples of each class')\ndf_t.groupBy('Survived', 'Sex').count().show()\n\n#Group By\nprint('---- Example_3 ---- Grouping by Survived and get the age mean (avg)')\ndf_t.groupBy('Survived').agg(F.mean(F.col('Age'))).show()\n\n\n#Group By\nprint('---- Example_3.1 ---- Grouping by Survived and get the age mean (avg) and Std. Also(Changing the column´s name)')\ndf_t.groupBy('Survived').agg(F.mean(F.col('Age')).alias('Mean Age'),\n                             F.stddev(F.col('Age')).alias('Std Age')).show()\n\n#Group By\nprint('---- Example_4 ---- Grouping by quality')\ndf.groupBy('quality').agg(F.mean(F.col('fixed acidity')).alias('fixed acidity mean'),\n                         F.stddev(F.col('fixed acidity')).alias('Std fixed acidity'),\n                         F.mean(F.col('residual sugar')).alias('residual sugar mean'),\n                         F.stddev(F.col('residual sugar')).alias('Std residual sugar')).show()\n\n#Group By\nprint('---- Example_5 ---- Grouping  age and counting only age > 20')\ndf_t.groupBy('Age').agg(F.when(F.col('Age') > 20, \\\n                            F.count(F.col('Age'))).alias('Count Age > 20')).show()","metadata":{"execution":{"iopub.status.busy":"2022-06-02T18:07:56.860315Z","iopub.execute_input":"2022-06-02T18:07:56.860649Z","iopub.status.idle":"2022-06-02T18:07:59.352424Z","shell.execute_reply.started":"2022-06-02T18:07:56.860604Z","shell.execute_reply":"2022-06-02T18:07:59.351481Z"},"trusted":true},"execution_count":null,"outputs":[]},{"cell_type":"markdown","source":"<a id=\"8\"></a>\n# <p style=\"background-color:#281F2F;height: 60px;text-align: center;vertical-align: middle;line-height: 60px;;font-family:helvetica;color:#FFFFFF;font-size:120%;text-align:center;border-radius:12px 12px;\"> WINDOW FUNCTION</p>","metadata":{}},{"cell_type":"code","source":"from pyspark.sql.window import Window\n\n#Creating a sequencial column in the time window\nprint('---- Example_1 ---- Creating a sequential Row')\nwindow  = Window.partitionBy(\"quality\").orderBy(\"quality\")\n\ndf_temp = df.withColumn(\"ROW COLUMN\",F.row_number().over(window))\ndf_temp.select('quality','alcohol','ROW COLUMN').show()\n\n#suming alcohol values over the quality variable\nprint('---- Example_2 ---- Sum over sequence')\nwindow  = Window.partitionBy(\"quality\").orderBy(\"quality\")\n\ndf_temp = df.withColumn(\"sum\",F.sum(F.col('alcohol')).over(window))\ndf_temp.select('quality','alcohol','sum').show()\n\n\n#Cumulative sum over alcohol by quality\nprint('---- Example_3 ---- Cumulative Sum over alcohol with ROUND')\nwindow  = Window.partitionBy(\"quality\").orderBy(\"alcohol\").rangeBetween(Window.unboundedPreceding, 0)\n\ndf_temp = df.withColumn(\"CUMULATIVE_alcohol\",F.round(F.sum(F.col('alcohol')).over(window),1))\ndf_temp.select('quality','alcohol','CUMULATIVE_alcohol').show()\n\n#Laging\nprint('---- Example_5 ---- Laging alcohol varibale by 2 (moving down 2 positions)')\nwindow  = Window.partitionBy(\"quality\").orderBy(\"alcohol\")\n\ndf_temp = df.withColumn(\"Lagged_alcohol\",F.lag(F.col('alcohol'),2).over(window))\ndf_temp.select('quality','alcohol','Lagged_alcohol').show()\n\n#Lead\nprint('---- Example_6 ---- Lead alcohol varibale by 2 (moving up 2 positions)')\nwindow  = Window.partitionBy(\"quality\").orderBy(\"alcohol\")\n\ndf_temp = df.withColumn(\"leaded_alcohol\",F.lead(F.col('alcohol'),2).over(window))\ndf_temp.select('quality','alcohol','leaded_alcohol').show()","metadata":{"execution":{"iopub.status.busy":"2022-06-02T18:07:59.353638Z","iopub.execute_input":"2022-06-02T18:07:59.354814Z","iopub.status.idle":"2022-06-02T18:08:01.322643Z","shell.execute_reply.started":"2022-06-02T18:07:59.354765Z","shell.execute_reply":"2022-06-02T18:08:01.321678Z"},"trusted":true},"execution_count":null,"outputs":[]},{"cell_type":"markdown","source":"<a id=\"9\"></a>\n# <p style=\"background-color:#281F2F;height: 60px;text-align: center;vertical-align: middle;line-height: 60px;;font-family:helvetica;color:#FFFFFF;font-size:120%;text-align:center;border-radius:12px 12px;\"> JOIN</p>\n\n<li style = \"font-size:110%\"> Join\n<li style = \"font-size:110%\"> unionByName\n<li style = \"font-size:110%\"> union \n\n","metadata":{}},{"cell_type":"code","source":"print('---- Example_1 ---- Join with a inner (If the keys don´t match,the value is droped)')\ndf_merged = df_example.join(df_example_2, on = ['key'], how = 'inner')\ndf_merged.show()\n\nprint('---- Example_2 ---- Join with a inner (Keep values from the left size of join (df_example) otherwise will become null values)')\ndf_merged = df_example.join(df_example_2, on = ['key'], how = 'left')\ndf_merged.show()\n\nprint('---- Example_3 ----Join with a inner (Keep values from the right size of join (df_example_2) otherwise will become null values)')\ndf_merged = df_example.join(df_example_2, on = ['key'], how = 'right')\ndf_merged.show()","metadata":{"execution":{"iopub.status.busy":"2022-06-02T18:40:38.162074Z","iopub.execute_input":"2022-06-02T18:40:38.162429Z","iopub.status.idle":"2022-06-02T18:40:39.368336Z","shell.execute_reply.started":"2022-06-02T18:40:38.162389Z","shell.execute_reply":"2022-06-02T18:40:39.36752Z"},"trusted":true},"execution_count":null,"outputs":[]},{"cell_type":"code","source":"#Creating a auxiliary df :)\nprint('Creating an auxiliary Df based on df_example_2 with functions learned above')\ndf_example_3 = df_example_2.withColumn('key', F.col('key')+F.lit(np.random.randint(50,100)))\ndf_example_3 = df_example_3.withColumn('C5', F.concat(F.col('C5'),F.lit('_F3')))\ndf_example_3 = df_example_3.withColumn('C6', F.round(F.col('C6') * F.lit(1.2)))\ndf_example_3.show()","metadata":{"execution":{"iopub.status.busy":"2022-06-02T19:01:49.640533Z","iopub.execute_input":"2022-06-02T19:01:49.640839Z","iopub.status.idle":"2022-06-02T19:01:49.850843Z","shell.execute_reply.started":"2022-06-02T19:01:49.640804Z","shell.execute_reply":"2022-06-02T19:01:49.850136Z"},"trusted":true},"execution_count":null,"outputs":[]},{"cell_type":"code","source":"print('---- Example_1 ---- Concatenating Rows by Name ( same column Names)')\ndf_union = df_example_2.unionByName(df_example_3)\ndf_union.show()","metadata":{"execution":{"iopub.status.busy":"2022-06-02T19:02:13.198097Z","iopub.execute_input":"2022-06-02T19:02:13.198435Z","iopub.status.idle":"2022-06-02T19:02:13.473503Z","shell.execute_reply.started":"2022-06-02T19:02:13.1984Z","shell.execute_reply":"2022-06-02T19:02:13.472545Z"},"trusted":true},"execution_count":null,"outputs":[]},{"cell_type":"code","source":"print('---- Example_2 ---- Concatenating Rows ( Atention to schema of columns in both dataset, should be the same)')\ndf_union = df_example_2.union(df_example_3)\ndf_union.show()","metadata":{"execution":{"iopub.status.busy":"2022-06-02T19:12:03.148514Z","iopub.execute_input":"2022-06-02T19:12:03.148833Z","iopub.status.idle":"2022-06-02T19:12:03.355256Z","shell.execute_reply.started":"2022-06-02T19:12:03.148794Z","shell.execute_reply":"2022-06-02T19:12:03.354274Z"},"trusted":true},"execution_count":null,"outputs":[]},{"cell_type":"markdown","source":"<a id=\"10\"></a>\n# <p style=\"background-color:#281F2F;height: 60px;text-align: center;vertical-align: middle;line-height: 60px;;font-family:helvetica;color:#FFFFFF;font-size:120%;text-align:center;border-radius:12px 12px;\"> UDF</p>\n\n<p style = 'font-size:120%'>User Defined Functions or UDF are a SQL wrapper of python functions. Hence, it is possible to create python functions and apply to transformer columns of a pyspark df </p>","metadata":{}},{"cell_type":"code","source":"\n# Let's say that we want to remove the F3 of column C5 of df_union previously created\n# First, create a dummy function. \n\ndef remove_F3(x):\n    \"\"\"Function to split string and return first element\"\"\"\n    return x.split('_')[0]\n\nprint(f'Teste String ABC_F3 -> {remove_F3(\"ABC_F3\")}\\n')\n\n# Next wrapping the functions into a udf(). The second arg is the type of the return element\nremove_F3_udf = F.udf(lambda x: remove_F3(x),StringType())\n\n#Finally, transform the C5 column\nprint('--- Example_1 --- Removing the F3 string from C5 Column with a custom function')\ndf_union_t = df_union.withColumn('C5_tranformed', remove_F3_udf(F.col('C5')))\ndf_union_t.show()","metadata":{"execution":{"iopub.status.busy":"2022-06-02T19:37:24.628313Z","iopub.execute_input":"2022-06-02T19:37:24.62863Z","iopub.status.idle":"2022-06-02T19:37:24.986219Z","shell.execute_reply.started":"2022-06-02T19:37:24.628593Z","shell.execute_reply":"2022-06-02T19:37:24.985562Z"},"trusted":true},"execution_count":null,"outputs":[]},{"cell_type":"code","source":"\n# Let's say that we want to combine the newly C5_tranformed (string) with a C6 (Number) multiplied by a factor (Why not?)\n# First, create a dummy function. \n\ndef combination_Func(st, num, factor):\n    \"\"\"Combine a string (st) with a num times a factor\"\"\"\n    return st+'_'+str(round(num*factor))\n\nprint(f'Teste String ABC \\ Number 10 \\ Factor 0.2 -> {combination_Func(\"ABC_F3\",10,0.2)}')\n\n# Next wrapping the functions into a udf(). The second arg is the type of the return element\ncombination_Func_udf = F.udf(lambda x,y,z: combination_Func(x,y,z),StringType())\n\n#Finally, transform the C5 column\nprint('--- Example_2 --- Creating a new column C7 based on C5_tranformed and C6 and function combination_Func')\ndf_union_t.withColumn('C7', combination_Func_udf(F.col('C5_tranformed'),F.col('C6'), F.lit(0.2))).show()","metadata":{"execution":{"iopub.status.busy":"2022-06-02T19:54:34.692249Z","iopub.execute_input":"2022-06-02T19:54:34.692566Z","iopub.status.idle":"2022-06-02T19:54:35.238815Z","shell.execute_reply.started":"2022-06-02T19:54:34.692534Z","shell.execute_reply":"2022-06-02T19:54:35.237878Z"},"trusted":true},"execution_count":null,"outputs":[]},{"cell_type":"markdown","source":"<a id=\"20\"></a>\n# <p style=\"background-color:#281F2F;height: 60px;text-align: center;vertical-align: middle;line-height: 60px;;font-family:helvetica;color:#FFFFFF;font-size:120%;text-align:center;border-radius:12px 12px;\"> REFERENCES</p>\n\n* [Spark Apache](https://spark.apache.org/)\n* [Spark Documentation](https://spark.apache.org/docs/2.4.0/api/python/)\n* [Spark SQL Functions](https://spark.apache.org/docs/2.4.0/api/python/pyspark.sql.html#pyspark.sql.functions)\n* (Book) Spark: The Definitive Guide: Big Data Processing Made Simple\n* (Book) Essential PySpark for Scalable Data Analytics: A beginner's guide to harnessing the power and ease of PySpark 3","metadata":{}},{"cell_type":"code","source":"","metadata":{},"execution_count":null,"outputs":[]},{"cell_type":"code","source":"","metadata":{},"execution_count":null,"outputs":[]}]}