{"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":"# **Using Observed Probability of the events as ratings**\n\n# Hypothesis and assumption:\n\n* All the purchases follows defined process: \n1. Users first clicks on the product for details\n1. Users adds the product to carts\n1. Users orders the product\n\n* As the data we are getting is for clicks and futher acitivities, we assume that the users would definitely click once they visit the website.\n\n> E.g.: for **session_clicks** for any particular product/aid : **ratings = p(aid|click)**,\n\n\n> Similarly for **session_carts** for any particular product/aid : **ratings = p(aid|carts|clicks)**  *{This is because of the instances observed where products are already in the Cart}*\n\n\n> And for **session_orders** for any particular product/aid : **ratings = p(aid|orders|carts|clicks)**\n\n\n* In this notebook- for any given item rating = **(cumulative distribution at that point)*(number of times the item triggered the event)/(total number of the particular events in the session)**\n\n* While it appears to ignore the probability of clicks converting to 'carts'/'orders', this could be included in future to improve performance\n\n* The timestamp is only used to get the number of clicks/carts/orders and not breakdown to time of the day/month/year","metadata":{}},{"cell_type":"code","source":"# This Python 3 environment comes with many helpful analytics libraries installed\n# It is defined by the kaggle/python Docker image: https://github.com/kaggle/docker-python\n# For example, here's several helpful packages to load\n\nimport numpy as np # linear algebra\nimport pandas as pd # data processing, CSV file I/O (e.g. pd.read_csv)\n\n# Input data files are available in the read-only \"../input/\" directory\n# For example, running this (by clicking run or pressing Shift+Enter) will list all files under the input directory\n\nimport os\nfor dirname, _, filenames in os.walk('/kaggle/input'):\n    for filename in filenames:\n        print(os.path.join(dirname, filename))\n\n# You can write up to 20GB to the current directory (/kaggle/working/) that gets preserved as output when you create a version using \"Save & Run All\" \n# You can also write temporary files to /kaggle/temp/, but they won't be saved outside of the current session","metadata":{"execution":{"iopub.status.busy":"2022-12-15T02:53:40.069032Z","iopub.execute_input":"2022-12-15T02:53:40.069823Z","iopub.status.idle":"2022-12-15T02:53:40.077436Z","shell.execute_reply.started":"2022-12-15T02:53:40.069783Z","shell.execute_reply":"2022-12-15T02:53:40.076552Z"},"trusted":true},"execution_count":null,"outputs":[]},{"cell_type":"code","source":"sample_sibmission = \"/kaggle/input/otto-recommender-system/sample_submission.csv\"\ntest = \"/kaggle/input/otto-recommender-system/test.jsonl\"\ntrain = \"/kaggle/input/otto-recommender-system/train.jsonl\"\n# df_train = pd.read_parquet(otto-chunk-data-inparquet-format)","metadata":{"execution":{"iopub.status.busy":"2022-12-15T02:53:42.601268Z","iopub.execute_input":"2022-12-15T02:53:42.602207Z","iopub.status.idle":"2022-12-15T02:53:42.606939Z","shell.execute_reply.started":"2022-12-15T02:53:42.602166Z","shell.execute_reply":"2022-12-15T02:53:42.605954Z"},"trusted":true},"execution_count":null,"outputs":[]},{"cell_type":"code","source":"%%time\n## Installing Apache Spark\n!pip install pyspark --quiet\n# Wall time: 52.6 s","metadata":{"execution":{"iopub.status.busy":"2022-12-15T02:53:44.638536Z","iopub.execute_input":"2022-12-15T02:53:44.639289Z","iopub.status.idle":"2022-12-15T02:54:30.567206Z","shell.execute_reply.started":"2022-12-15T02:53:44.639253Z","shell.execute_reply":"2022-12-15T02:54:30.565461Z"},"trusted":true},"execution_count":null,"outputs":[]},{"cell_type":"code","source":"# %%time\n# ## Installing polars\n# !pip install polars\n# #Wall time: 15.4 s","metadata":{"execution":{"iopub.status.busy":"2022-12-15T02:54:34.755930Z","iopub.execute_input":"2022-12-15T02:54:34.756368Z","iopub.status.idle":"2022-12-15T02:54:34.762870Z","shell.execute_reply.started":"2022-12-15T02:54:34.756328Z","shell.execute_reply":"2022-12-15T02:54:34.761554Z"},"trusted":true},"execution_count":null,"outputs":[]},{"cell_type":"code","source":"#Generic Libraries\nimport numpy as np # linear algebra\nimport pandas as pd # data processing, CSV file I/O (e.g. pd.read_csv)\n\n#Apache Spark Libraries\nimport pyspark\nfrom pyspark.context import SparkContext\nfrom pyspark.sql import SparkSession\n\n#Apache Spark SQL Functions\nfrom pyspark.sql import functions as F\nfrom pyspark.sql.window import Window\nfrom pyspark.sql.types import IntegerType\nfrom pyspark.sql.functions import explode, first, col, monotonically_increasing_id, lit, when, cume_dist, count, sum\n\n#Apache Spark ML CLassifier Libraries\nfrom pyspark.ml.classification import DecisionTreeClassifier,RandomForestClassifier,NaiveBayes\n\n#Apache Spark Evaluation Library\nfrom pyspark.ml.evaluation import MulticlassClassificationEvaluator\n\n#Apache Spark Features libraries\nfrom pyspark.ml.feature import StandardScaler,StringIndexer, VectorAssembler, VectorIndexer, OneHotEncoder, Normalizer, MinMaxScaler\n\n#Apache Spark Pipelin Library\nfrom pyspark.ml import Pipeline\n\n# Apache Spark `DenseVector`\nfrom pyspark.ml.linalg import DenseVector\n\n#Data Split Libraries\nimport sklearn\nfrom sklearn.model_selection import train_test_split\n\n#Polars to read files quickly\n\n\n# Import the requisite packages\nfrom pyspark.ml.tuning import ParamGridBuilder, CrossValidator\nfrom pyspark.ml.evaluation import RegressionEvaluator\n\n\n\n#Tabulating Data\nfrom tabulate import tabulate\n\n#Garbage\nimport gc","metadata":{"execution":{"iopub.status.busy":"2022-12-15T02:54:35.942847Z","iopub.execute_input":"2022-12-15T02:54:35.943229Z","iopub.status.idle":"2022-12-15T02:54:37.109799Z","shell.execute_reply.started":"2022-12-15T02:54:35.943199Z","shell.execute_reply":"2022-12-15T02:54:37.108752Z"},"trusted":true},"execution_count":null,"outputs":[]},{"cell_type":"code","source":"gc.collect()\n","metadata":{"execution":{"iopub.status.busy":"2022-12-15T02:54:38.065884Z","iopub.execute_input":"2022-12-15T02:54:38.066519Z","iopub.status.idle":"2022-12-15T02:54:38.209358Z","shell.execute_reply.started":"2022-12-15T02:54:38.066476Z","shell.execute_reply":"2022-12-15T02:54:38.208037Z"},"trusted":true},"execution_count":null,"outputs":[]},{"cell_type":"code","source":"import multiprocessing\ncores = multiprocessing.cpu_count()\nprint(cores)","metadata":{"execution":{"iopub.status.busy":"2022-12-15T02:54:39.162843Z","iopub.execute_input":"2022-12-15T02:54:39.163265Z","iopub.status.idle":"2022-12-15T02:54:39.169852Z","shell.execute_reply.started":"2022-12-15T02:54:39.163234Z","shell.execute_reply":"2022-12-15T02:54:39.168529Z"},"trusted":true},"execution_count":null,"outputs":[]},{"cell_type":"code","source":"spark = SparkSession.builder.\\\n                appName(\"FirstSparkApplication\").\\\n                config (\"spark.executor.memory\", \"25g\").\\\n                config (\"spark.default.parallelism\",\"400\").\\\n                config (\"spark.default.partitions\",\"10000\").\\\n                config (\"spark.sql.inMemoryColumnarStorage.compressed\", True).\\\n                config (\"spark.sql.inMemoryColumnarStorage.batchsize\", 10000).\\\n                getOrCreate()\nspark.sparkContext.setLogLevel('WARN')\nspark.sparkContext.version\nspark.sparkContext.setCheckpointDir(\"/kaggle/temp/\")","metadata":{"execution":{"iopub.status.busy":"2022-12-15T02:54:40.123215Z","iopub.execute_input":"2022-12-15T02:54:40.123923Z","iopub.status.idle":"2022-12-15T02:54:45.935613Z","shell.execute_reply.started":"2022-12-15T02:54:40.123884Z","shell.execute_reply":"2022-12-15T02:54:45.933855Z"},"trusted":true},"execution_count":null,"outputs":[]},{"cell_type":"code","source":"gc.collect()","metadata":{"execution":{"iopub.status.busy":"2022-12-15T02:54:50.576148Z","iopub.execute_input":"2022-12-15T02:54:50.576838Z","iopub.status.idle":"2022-12-15T02:54:50.731124Z","shell.execute_reply.started":"2022-12-15T02:54:50.576796Z","shell.execute_reply":"2022-12-15T02:54:50.729765Z"},"trusted":true},"execution_count":null,"outputs":[]},{"cell_type":"code","source":"%%time\n\"\"\"Defining Model\n\"\"\"\n\nfrom pyspark.ml.evaluation import RegressionEvaluator\nfrom pyspark.ml.recommendation import ALS, ALSModel\nfrom pyspark.ml.tuning import ParamGridBuilder, CrossValidator\nfrom pyspark.sql import Row\n\nals = ALS(rank=20, #10 was by default\n          maxIter=2 , regParam=0.01, \n          userCol=\"session_type\", itemCol=\"aid\", ratingCol=\"interest\", #ratings/possible_interest\n          coldStartStrategy=\"drop\",\n          implicitPrefs=True)\n#Wall time: 151 ms","metadata":{"execution":{"iopub.status.busy":"2022-12-15T02:54:59.945247Z","iopub.execute_input":"2022-12-15T02:54:59.945803Z","iopub.status.idle":"2022-12-15T02:55:01.834728Z","shell.execute_reply.started":"2022-12-15T02:54:59.945767Z","shell.execute_reply":"2022-12-15T02:55:01.833507Z"},"trusted":true},"execution_count":null,"outputs":[]},{"cell_type":"code","source":"gc.collect()","metadata":{"execution":{"iopub.status.busy":"2022-12-15T02:55:04.161965Z","iopub.execute_input":"2022-12-15T02:55:04.162397Z","iopub.status.idle":"2022-12-15T02:55:04.296660Z","shell.execute_reply.started":"2022-12-15T02:55:04.162366Z","shell.execute_reply":"2022-12-15T02:55:04.295445Z"},"trusted":true},"execution_count":null,"outputs":[]},{"cell_type":"code","source":"%%time\n\"\"\"importing and transforming Training Data\n\"\"\"\nprod_df_1 = spark.read.json(train)\ndf_1 = prod_df_1.na.drop() #Wall time: 2min 15s\ndf_sel = df_1.select(\"session\", explode(\"events\").alias(\"events\")).select(\"events.*\",\"*\").sort(\"ts\").select(\"aid\",\"ts\",\"session\",\"type\")\n\ndf_indexed = df_sel.withColumn(\"typeindex\", when(df_sel.type == \"clicks\",0)\n                                 .when(df_sel.type == \"carts\",1)\n                                 .otherwise(2))\n\ndf_session_type = df_indexed.withColumn(\"session_type\",(df_indexed.session)*10 + df_indexed.typeindex).select(\"aid\",\"ts\",\"session\",\"session_type\",\"type\")\n\n## Aggregating using \"Window Partition\"\nwindowPartition = Window.partitionBy(\"session_type\").orderBy(\"aid\")\n\ndf_freq=df_session_type.withColumn(\"cume_dist\",\n              cume_dist().over(windowPartition)).withColumn(\"frequency\",count(col(\"ts\")).over(windowPartition))\n\ndf_final = df_freq.withColumn(\"sumOfEvents\",sum(col(\"frequency\")).over(windowPartition))\n\ndf_train = df_final.select(\"aid\",\"session_type\",(df_final.cume_dist*(df_final.frequency / df_final.sumOfEvents)).alias('interest'))\n\ndf_train_check = df_train.checkpoint()\ndf_train_check.write.parquet(\"train.parquet\") \n# parDF1=spark.read.parquet(\"/temp/out/people.parquet\")\n\n#Wall time: 9min 41s with parallelism = 400/patition =4/executor_memory = 2G\n#Wall time: 13min 33s with parallelism = 400/patition =10000/executor_memory = 25G","metadata":{"execution":{"iopub.status.busy":"2022-12-15T02:55:06.564410Z","iopub.execute_input":"2022-12-15T02:55:06.564884Z","iopub.status.idle":"2022-12-15T03:10:19.837015Z","shell.execute_reply.started":"2022-12-15T02:55:06.564847Z","shell.execute_reply":"2022-12-15T03:10:19.831462Z"},"trusted":true},"execution_count":null,"outputs":[]},{"cell_type":"code","source":"gc.collect()","metadata":{"execution":{"iopub.status.busy":"2022-12-15T03:12:12.468166Z","iopub.execute_input":"2022-12-15T03:12:12.468919Z","iopub.status.idle":"2022-12-15T03:12:13.074825Z","shell.execute_reply.started":"2022-12-15T03:12:12.468863Z","shell.execute_reply":"2022-12-15T03:12:13.073356Z"},"trusted":true},"execution_count":null,"outputs":[]},{"cell_type":"code","source":"# %%time\n# model = als.fit(df_train_check)","metadata":{"execution":{"iopub.status.busy":"2022-12-15T03:12:16.695050Z","iopub.execute_input":"2022-12-15T03:12:16.695476Z","iopub.status.idle":"2022-12-15T03:12:16.699978Z","shell.execute_reply.started":"2022-12-15T03:12:16.695441Z","shell.execute_reply":"2022-12-15T03:12:16.699053Z"},"trusted":true},"execution_count":null,"outputs":[]},{"cell_type":"code","source":"%%time\n\"\"\"Importing and transforming data from Test File\n\"\"\"\n\nprod_df_test = spark.read.json(test)\ndf_test = prod_df_test.na.drop() #Wall time: 2min 15s\ndf_sel_test = df_test.select(\"session\", explode(\"events\").alias(\"events\")).select(\"events.*\",\"*\").sort(\"ts\").select(\"aid\",\"ts\",\"session\",\"type\")\ndf_indexed_test = df_sel_test.withColumn(\"typeindex\", when(df_sel_test.type == \"clicks\",0)\n                                 .when(df_sel_test.type == \"carts\",1)\n                                 .otherwise(2))\ndf_session_type_test = df_indexed_test.withColumn(\"session_type\",(df_indexed_test.session)*10 + df_indexed_test.typeindex).select(\"aid\",\"ts\",\"session\",\"session_type\",\"type\")\n## Aggregating using \"Window Partition\"\nwindowPartition = Window.partitionBy(\"session_type\").orderBy(\"aid\")\n\ndf_freq_test=df_session_type_test.withColumn(\"cume_dist\",\n              cume_dist().over(windowPartition)).withColumn(\"frequency\", \n                                                            count(col(\"ts\")).over(windowPartition))\n\ndf_final_test = df_freq_test.withColumn(\"sumOfEvents\",sum(col(\"frequency\")).over(windowPartition))\ndf_test = df_final_test.select(\"aid\",\"session_type\",(df_final_test.cume_dist*(df_final_test.frequency / df_final_test.sumOfEvents)).alias('interest')) #(df_final_test.cume_dist*(df_final_test.frequency / df_final_test.sumOfEvents)).alias('interest'))\ndf_test_check=df_test.checkpoint() #Wall time: 29.1 s #Wall time: 24.8 s\ndf_test_check.write.parquet(\"test.parquet\") ","metadata":{"execution":{"iopub.status.busy":"2022-12-15T03:12:19.874104Z","iopub.execute_input":"2022-12-15T03:12:19.875159Z","iopub.status.idle":"2022-12-15T03:12:53.690444Z","shell.execute_reply.started":"2022-12-15T03:12:19.875110Z","shell.execute_reply":"2022-12-15T03:12:53.689119Z"},"trusted":true},"execution_count":null,"outputs":[]},{"cell_type":"code","source":"# spark.sparkContext.stop()","metadata":{},"execution_count":null,"outputs":[]},{"cell_type":"code","source":"","metadata":{},"execution_count":null,"outputs":[]}]}