{"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":"code","source":"# Install pyspark\n!pip3 install pyspark","metadata":{"execution":{"iopub.status.busy":"2022-08-14T09:23:47.978080Z","iopub.execute_input":"2022-08-14T09:23:47.978559Z","iopub.status.idle":"2022-08-14T09:24:38.915912Z","shell.execute_reply.started":"2022-08-14T09:23:47.978456Z","shell.execute_reply":"2022-08-14T09:24:38.914155Z"},"trusted":true},"execution_count":null,"outputs":[]},{"cell_type":"markdown","source":"## Disclaimer :\nthis notebook is not published with good performance, the notebook aims to  propose an approach based on Spark to introduce people to Distributed Computing and to encourage this type of implementation, that is representative of the industry reality.\n\nIn this notebook, we'll try to introduce you  pyspark which is an interface for Apache Spark ( dedicated for distributed computation ) in Python.","metadata":{}},{"cell_type":"code","source":"import numpy as np\nimport pandas as pd \nimport random\nimport os","metadata":{"_uuid":"8f2839f25d086af736a60e9eeb907d3b93b6e0e5","_cell_guid":"b1076dfc-b9ad-4769-8c92-a6c4dae69d19","execution":{"iopub.status.busy":"2022-08-14T09:24:38.918905Z","iopub.execute_input":"2022-08-14T09:24:38.919339Z","iopub.status.idle":"2022-08-14T09:24:38.926337Z","shell.execute_reply.started":"2022-08-14T09:24:38.919299Z","shell.execute_reply":"2022-08-14T09:24:38.925144Z"},"trusted":true},"execution_count":null,"outputs":[]},{"cell_type":"markdown","source":"\n\nThe power of Pyspark is its simplicity, because we don't need to manipulate / use Resilient Distributed Dataset ( RDD ) which are the basic unit of Spark that are distributed on the cluster, pyspark handle all of this. However, if you want to program at this level you can easily do it.\n\n## Resilient Distributed Dataset\n-  An RDD (Resilient Distributed Dataset) is the basic abstraction of Spark representing an unchanging set of elements partitioned across cluster nodes, allowing parallel computation.\n- It has 3 advantages:\n    - Performance.\n    - Consistency.\n    - Fault Tolerance\n\nPyspark supports most of Spark features such as Dataframe / SparkSQL / MLlib. We will not introduce Spark Streaming in this tutorial.\n\n## SparkSQL \n- Simply said, sparkSQL allows  to query structured data inside Spark programs. It acts as a distributed SQL query engine for fast data retrieving and processing. We can  also use SQL's like code creating a temporary view as we will see at the end of the tutorial\n- It provides a programming abstraction called DataFrame:\n    -  Just as Python Dataframe, pyspark allows us to use this structure which is distributable across multiple machine of the cluster\n\n## MLlib\n- MLlib is a scalable machine learning library that provides a uniform set of high-level APIs that help users create and tune practical machine learning pipelines. For this competition, we'll use trees algorithms.\n\n","metadata":{}},{"cell_type":"code","source":"import pyspark\nimport pyspark.sql.functions as F\nimport pyspark.sql.types\nfrom pyspark.sql import SparkSession\nfrom pyspark.sql import SQLContext\nfrom pyspark import SparkContext\n\n\nfrom pyspark.ml.feature import PCA,StandardScaler,VectorAssembler\nfrom pyspark.ml import Pipeline\n\nfrom pyspark.ml.classification import DecisionTreeClassifier,RandomForestClassifier\nfrom pyspark.ml.feature import VectorAssembler\nfrom pyspark.ml.evaluation import BinaryClassificationEvaluator\nfrom sklearn.metrics import confusion_matrix","metadata":{"execution":{"iopub.status.busy":"2022-08-14T09:24:38.927584Z","iopub.execute_input":"2022-08-14T09:24:38.927874Z","iopub.status.idle":"2022-08-14T09:24:40.083144Z","shell.execute_reply.started":"2022-08-14T09:24:38.927848Z","shell.execute_reply":"2022-08-14T09:24:40.081939Z"},"trusted":true},"execution_count":null,"outputs":[]},{"cell_type":"code","source":"from pyspark.ml.feature import BucketedRandomProjectionLSH\nfrom pyspark.ml.linalg import Vectors, VectorUDT\nfrom pyspark.sql.functions import udf\nfrom pyspark.ml import Pipeline, Transformer\nfrom pyspark.sql import DataFrame\nfrom typing import Iterable\nimport pandas as pd","metadata":{"execution":{"iopub.status.busy":"2022-08-14T09:24:40.087378Z","iopub.execute_input":"2022-08-14T09:24:40.087789Z","iopub.status.idle":"2022-08-14T09:24:40.094166Z","shell.execute_reply.started":"2022-08-14T09:24:40.087746Z","shell.execute_reply":"2022-08-14T09:24:40.093023Z"},"trusted":true},"execution_count":null,"outputs":[]},{"cell_type":"markdown","source":"SparkSession is the entry point to Spark/Pyspark to work with RDD/Dataframe. We can also configure properties like :\n\n - `spark.driver.memory` : Amount of memory to use for the driver process where the Session is initialized. I initialize it by default at 16gb         \n - `spark.sql.shuffle.partitions` :  number of partitions to use when shuffling data for joins or aggregations. In this example we partitioned the data into 150 partitions.","metadata":{}},{"cell_type":"code","source":"def sample_df(n_cols,n_users,parDF1):\n    print(\"using PySpark ...\")\n    cols_list = [\"customer_ID\"] + random.sample(parDF1.columns[1:-1],n_cols) + ['target']\n    customerList = tuple(parDF1.select(\"customer_ID\").orderBy(F.rand()).limit(n_users).toPandas()['customer_ID'].values)\n    parDF1 = parDF1.select(*cols_list).filter(\"customer_ID IN {}\".format(customerList))\n    return parDF1","metadata":{"execution":{"iopub.status.busy":"2022-08-14T09:42:54.445122Z","iopub.execute_input":"2022-08-14T09:42:54.445562Z","iopub.status.idle":"2022-08-14T09:42:54.452167Z","shell.execute_reply.started":"2022-08-14T09:42:54.445501Z","shell.execute_reply":"2022-08-14T09:42:54.451355Z"},"trusted":true},"execution_count":null,"outputs":[]},{"cell_type":"code","source":"def agg_data(df):\n    \n    numerical_list = [item[0] for item in df.dtypes[1:] if not  item[1].startswith('string') ]\n\n    max_agg = {x: \"max\" for x in numerical_list   if x is not df.columns[0] }\n\n    max_df = df.groupBy(\"customer_ID\").agg(max_agg)\n\n    max_df.createOrReplaceTempView(\"max_df\")\n\n    sqlContext = SQLContext(spark)\n\n    all_temp_tables, all_temp_id = [], []\n    left_join_string = \"\"\n    for c in numerical_list[:-1]:\n        list_variables_random =         [F.avg(c).alias(f\"{c}_avg\"),\n                                         F.stddev_samp(c).alias(f\"{c}_std\"),\n                                         F.sum(c).alias(f\"{c}_sum\"),\n                                         F.skewness(c).alias(f\"{c}_skw\"),\n                                         F.kurtosis(c).alias(f\"{c}_kts\"),\n                                         F.min(c).alias(f\"{c}_min\")]\n                                         \n                    \n        \n        tmp = df.groupBy(\"customer_ID\").agg(*random.sample(list_variables_random,4)).drop(c).withColumnRenamed(\"customer_ID\",\"customer_ID_{}\".format(c))\n        tmp.createOrReplaceTempView(\"tmp_{}\".format(c))\n        all_temp_id.append(\"customer_ID_{}\".format(c))\n        all_temp_tables.append(\"tmp_{}\".format(c))\n        left_join_string  += \" left join tmp_{0} on max_df.customer_ID = tmp_{0}.customer_ID_{0}\".format(c)\n\n    new_df = sqlContext.sql(\"SELECT * FROM max_df \"+left_join_string).drop(*all_temp_id)\n    \n    for tmp_table in all_temp_tables:\n        spark.catalog.dropTempView(tmp_table)\n    return new_df","metadata":{"execution":{"iopub.status.busy":"2022-08-14T09:24:40.108475Z","iopub.execute_input":"2022-08-14T09:24:40.109262Z","iopub.status.idle":"2022-08-14T09:24:40.124404Z","shell.execute_reply.started":"2022-08-14T09:24:40.109217Z","shell.execute_reply":"2022-08-14T09:24:40.123579Z"},"trusted":true},"execution_count":null,"outputs":[]},{"cell_type":"code","source":"from pyspark.ml.param.shared import HasOutputCols, Param, Params\nfrom pyspark.ml.util import DefaultParamsReadable, DefaultParamsWritable\nclass ColumnTransformer(Transformer,DefaultParamsReadable, DefaultParamsWritable):\n    \"\"\"\n    A custom Transformer which drops all columns that have at least one of the\n    words from the banned_list in the name.\n    \"\"\"\n\n    def __init__(self, col):\n        super(ColumnTransformer, self).__init__()\n        self.col = col\n\n    def _transform(self, df: DataFrame) -> DataFrame:\n\n        df = df.select(self.col,\"target\")\n        list_to_vector_udf = udf(lambda l: Vectors.dense(l), VectorUDT())\n        df = df.withColumn(self.col, list_to_vector_udf(df[self.col]))\n        return df","metadata":{"execution":{"iopub.status.busy":"2022-08-14T09:24:40.125878Z","iopub.execute_input":"2022-08-14T09:24:40.127048Z","iopub.status.idle":"2022-08-14T09:24:40.140379Z","shell.execute_reply.started":"2022-08-14T09:24:40.127003Z","shell.execute_reply":"2022-08-14T09:24:40.139502Z"},"trusted":true},"execution_count":null,"outputs":[]},{"cell_type":"code","source":"spark = (SparkSession.builder.master(\"local[*]\")\n                    .config(\"spark.driver.memory\",\"16g\")\n                    .config(\"spark.sql.shuffle.partitions\",200)\n                    .getOrCreate())","metadata":{"execution":{"iopub.status.busy":"2022-08-14T09:24:40.142832Z","iopub.execute_input":"2022-08-14T09:24:40.143752Z","iopub.status.idle":"2022-08-14T09:24:45.643719Z","shell.execute_reply.started":"2022-08-14T09:24:40.143704Z","shell.execute_reply":"2022-08-14T09:24:45.642658Z"},"trusted":true},"execution_count":null,"outputs":[]},{"cell_type":"code","source":"spark.sparkContext.setLogLevel(\"WARN\")\n","metadata":{"execution":{"iopub.status.busy":"2022-08-14T09:24:45.647052Z","iopub.execute_input":"2022-08-14T09:24:45.648996Z","iopub.status.idle":"2022-08-14T09:24:45.656356Z","shell.execute_reply.started":"2022-08-14T09:24:45.648941Z","shell.execute_reply":"2022-08-14T09:24:45.655543Z"},"trusted":true},"execution_count":null,"outputs":[]},{"cell_type":"code","source":"parDF1=spark.read.parquet(\"../input/amex-parquet/train_data.parquet\")","metadata":{"execution":{"iopub.status.busy":"2022-08-14T09:42:49.021796Z","iopub.execute_input":"2022-08-14T09:42:49.022213Z","iopub.status.idle":"2022-08-14T09:42:49.124086Z","shell.execute_reply.started":"2022-08-14T09:42:49.022179Z","shell.execute_reply":"2022-08-14T09:42:49.123018Z"},"trusted":true},"execution_count":null,"outputs":[]},{"cell_type":"code","source":"for i in range(4):\n    cols_model = []\n\n\n    df = sample_df(35,10000,parDF1)\n    cols_used = df.columns[:-1]\n\n    df_aggregate_train = agg_data(df)\n\n    agg_df = df_aggregate_train.withColumnRenamed(\"max(target)\", \"target\")\n    input_cols = agg_df.select(\"*\").drop(\"customer_ID\",\"target\").columns\n    agg_df = agg_df.na.fill(value=0,subset=input_cols)\n    \n    \n    va = VectorAssembler(inputCols = input_cols, outputCol='features',handleInvalid = \"skip\")\n    standardScaler = StandardScaler(inputCol=\"features\",outputCol=\"features_\")\n    brp = PCA(k=15,inputCol=\"features_\",outputCol=\"lsh_{}\".format(i))\n    column_transformer = ColumnTransformer(\"lsh_{}\".format(i))\n    dtc = DecisionTreeClassifier(\n          maxDepth=8,\n          featuresCol=\"lsh_{}\".format(i),\n          labelCol=\"target\")\n\n\n    pipeline = Pipeline(stages=[va, \n                            standardScaler,\n                            brp,\n                            column_transformer,\n                            dtc]).fit(agg_df)\n\n    cols_model.append(cols_used)\n    pipeline.save(\"/kaggle/working/model_{}\".format(i))\n    with open(\"cols_model_{}.txt\".format(i), \"w\") as output:\n        output.write(str(cols_used))\n    spark.catalog.clearCache()\n    \n\n","metadata":{"execution":{"iopub.status.busy":"2022-08-14T09:43:24.286427Z","iopub.execute_input":"2022-08-14T09:43:24.287754Z","iopub.status.idle":"2022-08-14T09:51:07.780761Z","shell.execute_reply.started":"2022-08-14T09:43:24.287704Z","shell.execute_reply":"2022-08-14T09:51:07.778993Z"},"trusted":true},"execution_count":null,"outputs":[]},{"cell_type":"code","source":"","metadata":{"trusted":true},"execution_count":null,"outputs":[]}]}