{"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":"# Handle Large Dataset with Pyspark\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-16T04:18:14.984144Z","iopub.execute_input":"2022-06-16T04:18:14.985232Z","iopub.status.idle":"2022-06-16T04:19:04.336293Z","shell.execute_reply.started":"2022-06-16T04:18:14.985128Z","shell.execute_reply":"2022-06-16T04:19:04.334807Z"},"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-16T04:19:10.541876Z","iopub.execute_input":"2022-06-16T04:19:10.54243Z","iopub.status.idle":"2022-06-16T04:19:10.618725Z","shell.execute_reply.started":"2022-06-16T04:19:10.542381Z","shell.execute_reply":"2022-06-16T04:19:10.617827Z"},"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-16T04:19:14.168381Z","iopub.execute_input":"2022-06-16T04:19:14.168834Z","iopub.status.idle":"2022-06-16T04:19:20.370532Z","shell.execute_reply.started":"2022-06-16T04:19:14.168799Z","shell.execute_reply":"2022-06-16T04:19:20.369632Z"},"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-16T04:19:22.595896Z","iopub.execute_input":"2022-06-16T04:19:22.596309Z","iopub.status.idle":"2022-06-16T04:19:22.603038Z","shell.execute_reply.started":"2022-06-16T04:19:22.59626Z","shell.execute_reply":"2022-06-16T04:19:22.602132Z"},"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-16T04:19:27.33384Z","iopub.execute_input":"2022-06-16T04:19:27.334232Z","iopub.status.idle":"2022-06-16T04:19:27.419059Z","shell.execute_reply.started":"2022-06-16T04:19:27.334201Z","shell.execute_reply":"2022-06-16T04:19:27.417917Z"},"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-16T04:19:34.630079Z","iopub.execute_input":"2022-06-16T04:19:34.630493Z","iopub.status.idle":"2022-06-16T04:19:34.638953Z","shell.execute_reply.started":"2022-06-16T04:19:34.63046Z","shell.execute_reply":"2022-06-16T04:19:34.637978Z"},"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-16T04:19:38.478175Z","iopub.execute_input":"2022-06-16T04:19:38.478708Z","iopub.status.idle":"2022-06-16T04:19:38.484067Z","shell.execute_reply.started":"2022-06-16T04:19:38.478669Z","shell.execute_reply":"2022-06-16T04:19:38.483088Z"},"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-16T04:19:42.703861Z","iopub.execute_input":"2022-06-16T04:19:42.704272Z","iopub.status.idle":"2022-06-16T04:19:42.711539Z","shell.execute_reply.started":"2022-06-16T04:19:42.704236Z","shell.execute_reply":"2022-06-16T04:19:42.710638Z"},"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-16T04:19:46.891714Z","iopub.execute_input":"2022-06-16T04:19:46.892895Z","iopub.status.idle":"2022-06-16T04:19:46.898888Z","shell.execute_reply.started":"2022-06-16T04:19:46.892849Z","shell.execute_reply":"2022-06-16T04:19:46.897968Z"},"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-16T04:19:49.798202Z","iopub.execute_input":"2022-06-16T04:19:49.798626Z","iopub.status.idle":"2022-06-16T04:19:49.805652Z","shell.execute_reply.started":"2022-06-16T04:19:49.798592Z","shell.execute_reply":"2022-06-16T04:19:49.804781Z"},"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-16T04:19:53.919517Z","iopub.execute_input":"2022-06-16T04:19:53.919955Z","iopub.status.idle":"2022-06-16T04:19:53.931671Z","shell.execute_reply.started":"2022-06-16T04:19:53.919918Z","shell.execute_reply":"2022-06-16T04:19:53.93071Z"},"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)\n\n# train_psdf = spark.read.options(header=True, inferSchema=True).csv(train_path)\n# test_psdf = spark.read.options(header=True, inferSchema=True).csv(test_path)\n# label_psdf = spark.read.options(header=True, inferSchema=True).csv(label_path)","metadata":{"execution":{"iopub.status.busy":"2022-06-16T04:19:57.496243Z","iopub.execute_input":"2022-06-16T04:19:57.496701Z","iopub.status.idle":"2022-06-16T04:20:01.535615Z","shell.execute_reply.started":"2022-06-16T04:19:57.496663Z","shell.execute_reply":"2022-06-16T04:20:01.534553Z"},"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-16T03:46:14.954395Z","iopub.execute_input":"2022-06-16T03:46:14.95485Z","iopub.status.idle":"2022-06-16T03:46:15.014061Z","shell.execute_reply.started":"2022-06-16T03:46:14.954801Z","shell.execute_reply":"2022-06-16T03:46:15.012489Z"},"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-16T04:20:03.901246Z","iopub.execute_input":"2022-06-16T04:20:03.901978Z","iopub.status.idle":"2022-06-16T04:21:07.550628Z","shell.execute_reply.started":"2022-06-16T04:20:03.901932Z","shell.execute_reply":"2022-06-16T04:21:07.549674Z"},"trusted":true},"execution_count":null,"outputs":[]},{"cell_type":"code","source":"train_psdf.createOrReplaceTempView(\"train_psdf\")\n\nz = spark.sql(\"select * from train_psdf limit 50000\").toPandas()","metadata":{"execution":{"iopub.status.busy":"2022-06-16T04:21:19.306505Z","iopub.execute_input":"2022-06-16T04:21:19.306886Z","iopub.status.idle":"2022-06-16T04:21:37.989773Z","shell.execute_reply.started":"2022-06-16T04:21:19.306856Z","shell.execute_reply":"2022-06-16T04:21:37.988828Z"},"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":"z['S_2'].value_counts()","metadata":{"execution":{"iopub.status.busy":"2022-06-16T04:44:10.127304Z","iopub.execute_input":"2022-06-16T04:44:10.127721Z","iopub.status.idle":"2022-06-16T04:44:10.152154Z","shell.execute_reply.started":"2022-06-16T04:44:10.127687Z","shell.execute_reply":"2022-06-16T04:44:10.151265Z"},"trusted":true},"execution_count":null,"outputs":[]},{"cell_type":"code","source":"z","metadata":{"execution":{"iopub.status.busy":"2022-06-16T04:43:58.066367Z","iopub.execute_input":"2022-06-16T04:43:58.067261Z","iopub.status.idle":"2022-06-16T04:43:58.150564Z","shell.execute_reply.started":"2022-06-16T04:43:58.067221Z","shell.execute_reply":"2022-06-16T04:43:58.149451Z"},"trusted":true},"execution_count":null,"outputs":[]},{"cell_type":"code","source":"","metadata":{},"execution_count":null,"outputs":[]}]}