{"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":"# Import libraries","metadata":{}},{"cell_type":"code","source":"# !pip install pyspark","metadata":{"execution":{"iopub.status.busy":"2022-12-20T10:01:11.639118Z","iopub.execute_input":"2022-12-20T10:01:11.639556Z","iopub.status.idle":"2022-12-20T10:01:11.644424Z","shell.execute_reply.started":"2022-12-20T10:01:11.639520Z","shell.execute_reply":"2022-12-20T10:01:11.643311Z"},"trusted":true},"execution_count":null,"outputs":[]},{"cell_type":"code","source":"import numpy  as np # linear algebra\nimport pandas as pd # data processing, CSV file I/O (e.g. pd.read_csv)\n\nfrom pyspark.sql import SparkSession, functions as f, DataFrame","metadata":{"execution":{"iopub.status.busy":"2022-12-20T11:44:04.896434Z","iopub.execute_input":"2022-12-20T11:44:04.896953Z","iopub.status.idle":"2022-12-20T11:44:05.063391Z","shell.execute_reply.started":"2022-12-20T11:44:04.896857Z","shell.execute_reply":"2022-12-20T11:44:05.061885Z"},"trusted":true},"execution_count":null,"outputs":[]},{"cell_type":"markdown","source":"# Create Spark session","metadata":{}},{"cell_type":"code","source":"# Create a SparkSession\nspark = SparkSession.builder.getOrCreate()","metadata":{"execution":{"iopub.status.busy":"2022-12-20T11:44:05.763437Z","iopub.execute_input":"2022-12-20T11:44:05.764496Z","iopub.status.idle":"2022-12-20T11:44:12.366979Z","shell.execute_reply.started":"2022-12-20T11:44:05.764455Z","shell.execute_reply":"2022-12-20T11:44:12.366151Z"},"trusted":true},"execution_count":null,"outputs":[]},{"cell_type":"code","source":"spark.conf.set(\"spark.sql.session.timeZone\", \"UTC\")","metadata":{"execution":{"iopub.status.busy":"2022-12-20T11:44:12.368194Z","iopub.execute_input":"2022-12-20T11:44:12.368449Z","iopub.status.idle":"2022-12-20T11:44:13.506346Z","shell.execute_reply.started":"2022-12-20T11:44:12.368424Z","shell.execute_reply":"2022-12-20T11:44:13.505268Z"},"trusted":true},"execution_count":null,"outputs":[]},{"cell_type":"markdown","source":"# Import functions","metadata":{}},{"cell_type":"markdown","source":"## .. for data pre-processing","metadata":{}},{"cell_type":"code","source":"def f_load_data_into_pyspark_df(file_path: str, sample_size: int = None) -> DataFrame:\n    \"\"\"\n    Reads-in the `jsonl` file from `file_path` and converts it into a tabular PySpark dataframe format\n    \n    :param file_path: str Filepath string to the train or test set to read\n    :param sample_size: int The number of rows to read-in\n    \"\"\"\n    \n    # Read the file into a dataframe\n    if sample_size == None:\n        df = spark.read.json(file_path, multiLine=False)\n    else:\n        df = spark.read.json(file_path, multiLine=False).limit(sample_size)\n\n    # Extract the values of the \"events\" column into a new dataframe\n    extracted_values = df\\\n        .select(\"session\",f.explode(\"events\").alias(\"list\"))\\\n        .rdd\\\n        .flatMap(lambda x: [(x[0],x[1][0], x[1][1], x[1][2])])\n    \n    # Create a PySpark dataframe using the extracted values, and assign the column names as below:\n    new_df = spark.createDataFrame(extracted_values, [\"session\",\"aid\", \"ts\", \"type\"])\n\n    # 'ts' is in milliseconds. Convert to seconds by dividing by 1000. Then convert to datetime object.\n    new_df = new_df.withColumn(\"ts\", f.round(f.col('ts')/1000, 0))\n    new_df = new_df.withColumn(\"datetime\", f.from_unixtime(\"ts\", 'yyyy-MM-dd HH:mm:ss'))\n\n    # Select columns and sort by session and datetime\n    new_df = new_df\\\n        .select('session', 'datetime', 'aid', 'type')\\\n        .orderBy(['session', 'datetime'])\n\n    return new_df","metadata":{"execution":{"iopub.status.busy":"2022-12-20T11:44:31.662011Z","iopub.execute_input":"2022-12-20T11:44:31.662406Z","iopub.status.idle":"2022-12-20T11:44:31.674337Z","shell.execute_reply.started":"2022-12-20T11:44:31.662374Z","shell.execute_reply":"2022-12-20T11:44:31.672780Z"},"trusted":true},"execution_count":null,"outputs":[]},{"cell_type":"markdown","source":"# Extract pre-processed datasets as parquet","metadata":{}},{"cell_type":"code","source":"%%time\n\ntrain = f_load_data_into_pyspark_df('/kaggle/input/otto-recommender-system/train.jsonl')\n\ntrain.write.mode(\"overwrite\").parquet(\"/kaggle/working/train/\")","metadata":{"execution":{"iopub.status.busy":"2022-12-20T11:44:36.083781Z","iopub.execute_input":"2022-12-20T11:44:36.084177Z","iopub.status.idle":"2022-12-20T12:47:48.047155Z","shell.execute_reply.started":"2022-12-20T11:44:36.084142Z","shell.execute_reply":"2022-12-20T12:47:48.042568Z"},"trusted":true},"execution_count":null,"outputs":[]},{"cell_type":"code","source":"%%time\n\ntest = f_load_data_into_pyspark_df('/kaggle/input/otto-recommender-system/test.jsonl')\ntest.write.mode(\"overwrite\").parquet(\"/kaggle/working/test/\")","metadata":{"execution":{"iopub.status.busy":"2022-12-20T12:47:48.061671Z","iopub.execute_input":"2022-12-20T12:47:48.063884Z","iopub.status.idle":"2022-12-20T12:49:52.069764Z","shell.execute_reply.started":"2022-12-20T12:47:48.063790Z","shell.execute_reply":"2022-12-20T12:49:52.068477Z"},"trusted":true},"execution_count":null,"outputs":[]},{"cell_type":"markdown","source":"## Read-in","metadata":{}},{"cell_type":"code","source":"train_df = spark.read.parquet(\"/kaggle/working/train/\").orderBy(['session', 'datetime'])\n\nprint(f\"{train_df.count():,}, {len(train_df.columns)}\")\n\ntrain_df.limit(10).toPandas()","metadata":{"execution":{"iopub.status.busy":"2022-12-20T12:59:49.567113Z","iopub.execute_input":"2022-12-20T12:59:49.567578Z","iopub.status.idle":"2022-12-20T13:00:18.201762Z","shell.execute_reply.started":"2022-12-20T12:59:49.567530Z","shell.execute_reply":"2022-12-20T13:00:18.200571Z"},"trusted":true},"execution_count":null,"outputs":[]},{"cell_type":"code","source":"test_df = spark.read.parquet(\"/kaggle/working/test/\").orderBy(['session', 'datetime'])\n\nprint(f\"{test_df.count():,}, {len(test_df.columns)}\")\n\ntest_df.limit(10).toPandas()","metadata":{"execution":{"iopub.status.busy":"2022-12-20T13:00:56.576439Z","iopub.execute_input":"2022-12-20T13:00:56.576887Z","iopub.status.idle":"2022-12-20T13:00:57.711226Z","shell.execute_reply.started":"2022-12-20T13:00:56.576852Z","shell.execute_reply":"2022-12-20T13:00:57.710084Z"},"trusted":true},"execution_count":null,"outputs":[]},{"cell_type":"code","source":"","metadata":{},"execution_count":null,"outputs":[]}]}