{"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":"# Export Large Dataset to Spark\nThe current competition data file size is about 30GB+. We can't use `pandas` because it will take a long time to even read it. You can switch to `dask` or `pyspark` for this kind of problem. In this notebook I will choose `pyspark` because I want to apply what have I learned about data engineering.\n\nThe real reason we use `pyspark` is because it runs operations using multiple machine while `pandas` only use single machine. `pyspark` can perform lazy operation so that we have no to wait every operations to be finished. If you try to read the data using `pandas` it will take a long long time to even finished the read operations, that isn't a good practice. Some may use `cudf` and that is quite a good idea, but we will stick to `pyspark`.\n\nThe downside is `pyspark` have less algorithms than `pandas`, it might restraining our flexibility.","metadata":{}},{"cell_type":"markdown","source":"## Install pyspark","metadata":{}},{"cell_type":"code","source":"!pip install -q pyspark","metadata":{"execution":{"iopub.status.busy":"2022-06-16T06:15:39.139327Z","iopub.execute_input":"2022-06-16T06:15:39.140458Z","iopub.status.idle":"2022-06-16T06:16:27.222589Z","shell.execute_reply.started":"2022-06-16T06:15:39.14032Z","shell.execute_reply":"2022-06-16T06:16:27.221562Z"},"trusted":true},"execution_count":null,"outputs":[]},{"cell_type":"markdown","source":"## Import Libraries","metadata":{}},{"cell_type":"code","source":"import os\nfrom pprint import pprint\n\nimport pandas as pd\nfrom pyspark.sql import SparkSession, types","metadata":{"_uuid":"8f2839f25d086af736a60e9eeb907d3b93b6e0e5","_cell_guid":"b1076dfc-b9ad-4769-8c92-a6c4dae69d19","execution":{"iopub.status.busy":"2022-06-16T06:16:27.224879Z","iopub.execute_input":"2022-06-16T06:16:27.225509Z","iopub.status.idle":"2022-06-16T06:16:27.292298Z","shell.execute_reply.started":"2022-06-16T06:16:27.225451Z","shell.execute_reply":"2022-06-16T06:16:27.29135Z"},"trusted":true},"execution_count":null,"outputs":[]},{"cell_type":"markdown","source":"## Spark Session\nIn order to use `pyspark` we need to create or get `spark` instance.","metadata":{}},{"cell_type":"code","source":"spark = SparkSession.builder.master(\"local[*]\").getOrCreate()","metadata":{"execution":{"iopub.status.busy":"2022-06-16T06:16:27.293408Z","iopub.execute_input":"2022-06-16T06:16:27.293732Z","iopub.status.idle":"2022-06-16T06:16:32.449811Z","shell.execute_reply.started":"2022-06-16T06:16:27.29369Z","shell.execute_reply":"2022-06-16T06:16:32.44873Z"},"trusted":true},"execution_count":null,"outputs":[]},{"cell_type":"markdown","source":"## Infer Data Types\nWhen we read the data directly with `pyspark` it regards all the data types as string. While for `pandas`, it tries to infer what the data types by the value. We will use `pandas` for this task, but we can't afford to read all the data by whole. However, we only read the first n rows of the data and fetch all the data types.","metadata":{}},{"cell_type":"code","source":"# Define data paths\ntest_path = \"../input/amex-default-prediction/test_data.csv\"\ntrain_path = \"../input/amex-default-prediction/train_data.csv\"\nlabel_path = \"../input/amex-default-prediction/train_labels.csv\"","metadata":{"execution":{"iopub.status.busy":"2022-06-16T06:16:32.452166Z","iopub.execute_input":"2022-06-16T06:16:32.453973Z","iopub.status.idle":"2022-06-16T06:16:32.462932Z","shell.execute_reply.started":"2022-06-16T06:16:32.453917Z","shell.execute_reply":"2022-06-16T06:16:32.461112Z"},"trusted":true},"execution_count":null,"outputs":[]},{"cell_type":"code","source":"# Read data\ntrain_df = pd.read_csv(train_path, nrows=100)\ntest_df = pd.read_csv(test_path, nrows=100)\nlabel_df = pd.read_csv(label_path, nrows=100)","metadata":{"execution":{"iopub.status.busy":"2022-06-16T06:16:32.464178Z","iopub.execute_input":"2022-06-16T06:16:32.464754Z","iopub.status.idle":"2022-06-16T06:16:32.56243Z","shell.execute_reply.started":"2022-06-16T06:16:32.464719Z","shell.execute_reply":"2022-06-16T06:16:32.56174Z"},"trusted":true},"execution_count":null,"outputs":[]},{"cell_type":"code","source":"# Get all data types\n\n## Train types\ntrain_types = train_df.dtypes\ntrain_types_count = train_types.value_counts()\n\n## Test types\ntest_types = test_df.dtypes\ntest_types_count = test_types.value_counts()\n\n## Label types\nlabel_types = label_df.dtypes\nlabel_types_count = label_types.value_counts()","metadata":{"execution":{"iopub.status.busy":"2022-06-16T06:16:32.563352Z","iopub.execute_input":"2022-06-16T06:16:32.563925Z","iopub.status.idle":"2022-06-16T06:16:32.576153Z","shell.execute_reply.started":"2022-06-16T06:16:32.56389Z","shell.execute_reply":"2022-06-16T06:16:32.575485Z"},"trusted":true},"execution_count":null,"outputs":[]},{"cell_type":"code","source":"def print_splits(*msg):\n    for m in msg:\n        print(m)\n        print()","metadata":{"execution":{"iopub.status.busy":"2022-06-16T06:16:32.577855Z","iopub.execute_input":"2022-06-16T06:16:32.578712Z","iopub.status.idle":"2022-06-16T06:16:32.583957Z","shell.execute_reply.started":"2022-06-16T06:16:32.578648Z","shell.execute_reply":"2022-06-16T06:16:32.583053Z"},"trusted":true},"execution_count":null,"outputs":[]},{"cell_type":"markdown","source":"We get all the data types","metadata":{}},{"cell_type":"code","source":"print_splits(train_types_count, test_types_count, label_types_count)","metadata":{"execution":{"iopub.status.busy":"2022-06-16T06:16:32.585151Z","iopub.execute_input":"2022-06-16T06:16:32.586038Z","iopub.status.idle":"2022-06-16T06:16:32.594748Z","shell.execute_reply.started":"2022-06-16T06:16:32.585997Z","shell.execute_reply":"2022-06-16T06:16:32.593884Z"},"trusted":true},"execution_count":null,"outputs":[]},{"cell_type":"markdown","source":"## Create Schemas\nAfter we have retrieved the data types, we can create spark schema by creating `StructType` instances with every column as the argument. Every column is defined using `StructField` instance, it receives 3 arguments `colname`, `data types` and `nullable`.","metadata":{}},{"cell_type":"code","source":"# Types mapper\ntypes_map = {\n    \"object\": types.StringType(),\n    \"float64\": types.FloatType(),\n    \"int64\": types.IntegerType(),\n}\n\n# Known dtypes\nstring_dtypes = ['B_30', 'B_38', 'D_114', 'D_116', 'D_117', 'D_120', 'D_126', 'D_63', 'D_64', 'D_66', 'D_68']\ndate_dtypes = ['S_2']","metadata":{"execution":{"iopub.status.busy":"2022-06-16T06:16:32.596Z","iopub.execute_input":"2022-06-16T06:16:32.596468Z","iopub.status.idle":"2022-06-16T06:16:32.606376Z","shell.execute_reply.started":"2022-06-16T06:16:32.596439Z","shell.execute_reply":"2022-06-16T06:16:32.605719Z"},"trusted":true},"execution_count":null,"outputs":[]},{"cell_type":"code","source":"def create_spark_schema(series):\n    fields = []\n    \n    for index, value in series.items():\n        if index in string_dtypes:\n            field = types.StructField(index, types.StringType(), True)\n            \n        elif index in date_dtypes:\n            field = types.StructField(index, types.DateType(), True)\n        \n        else:\n            field = types.StructField(index, types_map.get(str(value)), True)\n            \n        fields.append(field)\n    return types.StructType(fields)","metadata":{"execution":{"iopub.status.busy":"2022-06-16T06:16:32.608068Z","iopub.execute_input":"2022-06-16T06:16:32.608665Z","iopub.status.idle":"2022-06-16T06:16:32.616385Z","shell.execute_reply.started":"2022-06-16T06:16:32.608632Z","shell.execute_reply":"2022-06-16T06:16:32.615745Z"},"trusted":true},"execution_count":null,"outputs":[]},{"cell_type":"code","source":"train_schema = create_spark_schema(train_types) \ntest_schema = create_spark_schema(test_types)\nlabel_schema = create_spark_schema(label_types)","metadata":{"execution":{"iopub.status.busy":"2022-06-16T06:16:32.617327Z","iopub.execute_input":"2022-06-16T06:16:32.618015Z","iopub.status.idle":"2022-06-16T06:16:32.633644Z","shell.execute_reply.started":"2022-06-16T06:16:32.617967Z","shell.execute_reply":"2022-06-16T06:16:32.632904Z"},"trusted":true},"execution_count":null,"outputs":[]},{"cell_type":"markdown","source":"## Read with Pyspark","metadata":{}},{"cell_type":"code","source":"# Set header to True or else it will be included as row\ntrain_psdf = spark.read.option(\"header\", \"true\").csv(train_path, schema=train_schema)\ntest_psdf = spark.read.option(\"header\", \"true\").csv(test_path, schema=test_schema)\nlabel_psdf = spark.read.option(\"header\", \"true\").csv(label_path, schema=label_schema)","metadata":{"execution":{"iopub.status.busy":"2022-06-16T06:16:32.634852Z","iopub.execute_input":"2022-06-16T06:16:32.635758Z","iopub.status.idle":"2022-06-16T06:16:36.257654Z","shell.execute_reply.started":"2022-06-16T06:16:32.635713Z","shell.execute_reply":"2022-06-16T06:16:36.25671Z"},"trusted":true},"execution_count":null,"outputs":[]},{"cell_type":"code","source":"# Check schema\nprint_splits(train_psdf.schema[:3], test_psdf.schema[:3], label_psdf.schema[:3])","metadata":{"execution":{"iopub.status.busy":"2022-06-16T06:00:18.294256Z","iopub.execute_input":"2022-06-16T06:00:18.29462Z","iopub.status.idle":"2022-06-16T06:00:18.345638Z","shell.execute_reply.started":"2022-06-16T06:00:18.294592Z","shell.execute_reply":"2022-06-16T06:00:18.344812Z"},"trusted":true},"execution_count":null,"outputs":[]},{"cell_type":"markdown","source":"## Counts Data\nThe data is about 5 millions rows with 100+ columns","metadata":{}},{"cell_type":"code","source":"train_psdf.count()","metadata":{"execution":{"iopub.status.busy":"2022-06-16T06:02:23.806609Z","iopub.execute_input":"2022-06-16T06:02:23.807445Z","iopub.status.idle":"2022-06-16T06:03:30.522161Z","shell.execute_reply.started":"2022-06-16T06:02:23.807412Z","shell.execute_reply":"2022-06-16T06:03:30.521393Z"},"trusted":true},"execution_count":null,"outputs":[]},{"cell_type":"markdown","source":"## Save as Parquet\nWe choose `.parquet` as the file extension because it uses less disk memory.","metadata":{}},{"cell_type":"code","source":"train_psdf.write.parquet(\"train_amex\")\ntest_psdf.write.parquet(\"test_amex\")\nlabel_psdf.write.parquet(\"label_amex\")","metadata":{"execution":{"iopub.status.busy":"2022-06-16T06:03:30.523477Z","iopub.execute_input":"2022-06-16T06:03:30.523845Z","iopub.status.idle":"2022-06-16T06:06:27.826481Z","shell.execute_reply.started":"2022-06-16T06:03:30.523793Z","shell.execute_reply":"2022-06-16T06:06:27.824845Z"},"collapsed":true,"jupyter":{"outputs_hidden":true},"trusted":true},"execution_count":null,"outputs":[]},{"cell_type":"markdown","source":"## What to do next?\n1. You can read the exported data with pyspark API it can be `pyspark.sql`, `pyspark.pandas`, or `pyspark.rdd`\n2. Then perform preprocesing, you can see the related notebook from [Michal Slapek](https://www.kaggle.com/capslock): [here](https://www.kaggle.com/code/capslock/amex-export-to-parquet-with-apache-spark)\n\nMy coverage notebook about doing ML in pyspark for this competitions is still in progress. If anybody have done it I would like to know :D how it is done.","metadata":{}},{"cell_type":"markdown","source":"## Update\n\n**Changes:**\n- v.4 Added `.option(\"header\", \"true\")` to not include header as row\n- v.5 Update schema with known dtypes\n- v.6 Failed Run\n- v.7 Change the title to avoid misleading (the preprocessing notebook is still in progress)\n\n\nGood luck for the competitions :D.","metadata":{}},{"cell_type":"code","source":"","metadata":{},"execution_count":null,"outputs":[]}]}