{"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", "metadata": {}, "source": "# Big Data Engineering: 50GB CSV to Parquet via PySpark\n\n**Portfolio Concept:** The American Express Default Prediction dataset is notoriously massive (~50GB of raw CSVs). Loading this into Pandas directly on a standard machine causes fatal Out-Of-Memory (OOM) errors.\n\nIn this notebook, we utilize **PySpark** to ingest the data in distributed chunks, optimize the data types mathematically, and rewrite the dataset into highly compressed, partitioned **Parquet** files. This reduces the footprint dramatically and enables rapid downstream modeling."}, {"cell_type": "markdown", "metadata": {}, "source": "# Big Data Engineering: 50GB CSV to Parquet via PySpark\n\n**Portfolio Concept:** The American Express Default Prediction dataset is notoriously massive (~50GB of raw CSVs). Loading this into Pandas directly on a standard machine causes fatal Out-Of-Memory (OOM) errors.\n\nIn this notebook, we utilize **PySpark** to ingest the data in distributed chunks, optimize the data types mathematically, and rewrite the dataset into highly compressed, partitioned **Parquet** files. This reduces the footprint dramatically and enables rapid downstream modeling."}, {"cell_type": "code", "source": "import numpy as np \nimport pandas as pd\n\n\"\"\"import shutil\nshutil.rmtree('/kaggle/working/amexparq')\"\"\"\n\nimport os\nfor dirname, _, filenames in os.walk('/kaggle/input'):\n    for filename in filenames:\n        print(os.path.join(dirname, filename))", "metadata": {"_uuid": "8f2839f25d086af736a60e9eeb907d3b93b6e0e5", "_cell_guid": "b1076dfc-b9ad-4769-8c92-a6c4dae69d19", "execution": {"iopub.status.busy": "2022-08-02T00:01:11.949536Z", "iopub.execute_input": "2022-08-02T00:01:11.95Z", "iopub.status.idle": "2022-08-02T00:01:11.984079Z", "shell.execute_reply.started": "2022-08-02T00:01:11.949913Z", "shell.execute_reply": "2022-08-02T00:01:11.982861Z"}, "trusted": true}, "execution_count": null, "outputs": []}, {"cell_type": "code", "source": "\"\"\"import shutil\nshutil.rmtree('/kaggle/working/amexparq')\"\"\"", "metadata": {"execution": {"iopub.status.busy": "2022-08-02T00:01:11.986172Z", "iopub.execute_input": "2022-08-02T00:01:11.987028Z", "iopub.status.idle": "2022-08-02T00:01:11.996344Z", "shell.execute_reply.started": "2022-08-02T00:01:11.986976Z", "shell.execute_reply": "2022-08-02T00:01:11.995017Z"}, "trusted": true}, "execution_count": null, "outputs": []}, {"cell_type": "code", "source": "! pip install pyspark", "metadata": {"execution": {"iopub.status.busy": "2022-08-02T00:01:11.998464Z", "iopub.execute_input": "2022-08-02T00:01:11.999435Z", "iopub.status.idle": "2022-08-02T00:02:21.441138Z", "shell.execute_reply.started": "2022-08-02T00:01:11.999386Z", "shell.execute_reply": "2022-08-02T00:02:21.439164Z"}, "trusted": true}, "execution_count": null, "outputs": []}, {"cell_type": "markdown", "metadata": {}, "source": "## 1. PySpark Initialization\nWe initialize a local Spark session, allocating 14GB of driver memory and setting shuffle partitions. This configures the JVM to handle heavy distributed processing on a single node."}, {"cell_type": "markdown", "metadata": {}, "source": "## 1. PySpark Initialization\nWe initialize a local Spark session, allocating 14GB of driver memory and setting shuffle partitions. This configures the JVM to handle heavy distributed processing on a single node."}, {"cell_type": "code", "source": "import pyspark\nimport pyspark.sql.functions as F\nimport pyspark.sql.types\nfrom pyspark.sql import SparkSession ", "metadata": {"execution": {"iopub.status.busy": "2022-08-02T00:02:21.444369Z", "iopub.execute_input": "2022-08-02T00:02:21.44492Z", "iopub.status.idle": "2022-08-02T00:02:21.534332Z", "shell.execute_reply.started": "2022-08-02T00:02:21.444879Z", "shell.execute_reply": "2022-08-02T00:02:21.532939Z"}, "trusted": true}, "execution_count": null, "outputs": []}, {"cell_type": "code", "source": "print(\"Pyspark Version\",pyspark.__version__)", "metadata": {"execution": {"iopub.status.busy": "2022-08-02T00:02:21.53574Z", "iopub.execute_input": "2022-08-02T00:02:21.536127Z", "iopub.status.idle": "2022-08-02T00:02:21.543013Z", "shell.execute_reply.started": "2022-08-02T00:02:21.536092Z", "shell.execute_reply": "2022-08-02T00:02:21.541711Z"}, "trusted": true}, "execution_count": null, "outputs": []}, {"cell_type": "code", "source": "spark = SparkSession.builder.master(\"local[*]\").config(\"spark.driver.memory\",\"14g\").getOrCreate()\nspark.conf.set(\"spark.sql.shuffle.partitions\",20)", "metadata": {"execution": {"iopub.status.busy": "2022-08-02T00:02:21.546799Z", "iopub.execute_input": "2022-08-02T00:02:21.547649Z", "iopub.status.idle": "2022-08-02T00:02:28.847382Z", "shell.execute_reply.started": "2022-08-02T00:02:21.547575Z", "shell.execute_reply": "2022-08-02T00:02:28.846208Z"}, "trusted": true}, "execution_count": null, "outputs": []}, {"cell_type": "code", "source": "# from Dataset description\ncat_col = {'B_30','B_38','D_114','D_116','D_117','D_120','D_126','D_66','D_68','D_63','D_64'}\ncat_lab_col = {'D_63','D_64'}", "metadata": {"execution": {"iopub.status.busy": "2022-08-02T00:02:28.848674Z", "iopub.execute_input": "2022-08-02T00:02:28.849117Z", "iopub.status.idle": "2022-08-02T00:02:28.860057Z", "shell.execute_reply.started": "2022-08-02T00:02:28.84908Z", "shell.execute_reply": "2022-08-02T00:02:28.858085Z"}, "trusted": true}, "execution_count": null, "outputs": []}, {"cell_type": "code", "source": "col_nm_list = ['customer_ID','S_2','P_2','D_39','B_1','B_2','R_1','S_3','D_41','B_3','D_42','D_43','D_44','B_4','D_45','B_5','R_2',\n 'D_46','D_47','D_48','D_49','B_6','B_7','B_8','D_50','D_51','B_9','R_3','D_52','P_3','B_10','D_53','S_5','B_11','S_6','D_54','R_4',\n 'S_7','B_12','S_8','D_55','D_56','B_13','R_5','D_58','S_9','B_14','D_59','D_60','D_61','B_15','S_11','D_62','D_63','D_64','D_65',\n 'B_16','B_17','B_18','B_19','D_66','B_20','D_68','S_12','R_6','S_13','B_21','D_69','B_22','D_70','D_71','D_72','S_15','B_23','D_73',\n 'P_4','D_74','D_75','D_76','B_24','R_7','D_77','B_25','B_26','D_78','D_79','R_8','R_9','S_16','D_80','R_10','R_11','B_27','D_81','D_82',\n 'S_17','R_12','B_28','R_13','D_83','R_14','R_15','D_84','R_16','B_29','B_30','S_18','D_86','D_87','R_17','R_18','D_88','B_31','S_19',\n 'R_19','B_32','S_20','R_20','R_21','B_33','D_89','R_22','R_23','D_91','D_92','D_93','D_94','R_24','R_25','D_96','S_22','S_23',\n 'S_24','S_25','S_26','D_102','D_103','D_104','D_105','D_106','D_107','B_36','B_37','R_26','R_27','B_38','D_108',\n 'D_109','D_110','D_111','B_39','D_112','B_40','S_27','D_113','D_114','D_115','D_116','D_117','D_118','D_119',\n 'D_120','D_121','D_122','D_123','D_124','D_125','D_126','D_127','D_128','D_129','B_41','B_42','D_130','D_131',\n 'D_132','D_133','R_28','D_134','D_135','D_136','D_137','D_138','D_139','D_140','D_141','D_142','D_143','D_144','D_145']", "metadata": {"execution": {"iopub.status.busy": "2022-08-02T00:02:28.86221Z", "iopub.execute_input": "2022-08-02T00:02:28.862774Z", "iopub.status.idle": "2022-08-02T00:02:28.88038Z", "shell.execute_reply.started": "2022-08-02T00:02:28.862715Z", "shell.execute_reply": "2022-08-02T00:02:28.879206Z"}, "trusted": true}, "execution_count": null, "outputs": []}, {"cell_type": "markdown", "metadata": {}, "source": "## 2. Hardcoded Schema Definition\n**Why hardcode the schema?** If we tell Spark to `inferSchema=True`, it has to read the entire 50GB CSV twice (once to figure out the column types, and once to actually load the data). By explicitly defining the `StructType` mapping beforehand, we cut the I/O load time in half!"}, {"cell_type": "markdown", "metadata": {}, "source": "## 2. Hardcoded Schema Definition\n**Why hardcode the schema?** If we tell Spark to `inferSchema=True`, it has to read the entire 50GB CSV twice (once to figure out the column types, and once to actually load the data). By explicitly defining the `StructType` mapping beforehand, we cut the I/O load time in half!"}, {"cell_type": "code", "source": "def col_dtype(col_nm):\n    if col_nm == 'customer_ID':\n        return pyspark.sql.types.StringType()\n    elif col_nm == 'S_2':\n        return pyspark.sql.types.DateType()\n    elif col_nm in cat_col:\n        if col_nm in cat_lab_col:\n            return pyspark.sql.types.StringType()\n        else:\n            return pyspark.sql.types.FloatType()\n    else:\n        return pyspark.sql.types.FloatType()\n    \nschema = pyspark.sql.types.StructType([pyspark.sql.types.StructField(col,col_dtype(col)) for col in col_nm_list])         ", "metadata": {"execution": {"iopub.status.busy": "2022-08-02T00:02:28.883359Z", "iopub.execute_input": "2022-08-02T00:02:28.884497Z", "iopub.status.idle": "2022-08-02T00:02:28.898511Z", "shell.execute_reply.started": "2022-08-02T00:02:28.884452Z", "shell.execute_reply": "2022-08-02T00:02:28.897261Z"}, "trusted": true}, "execution_count": null, "outputs": []}, {"cell_type": "markdown", "metadata": {}, "source": "## 3. Streaming the Data\nUsing `FAILFAST` mode ensures that if any row violates our strict schema, the job halts immediately rather than silently injecting nulls into our clean data pipeline."}, {"cell_type": "markdown", "metadata": {}, "source": "## 3. Streaming the Data\nUsing `FAILFAST` mode ensures that if any row violates our strict schema, the job halts immediately rather than silently injecting nulls into our clean data pipeline."}, {"cell_type": "code", "source": "def read_file(path):\n    limit = None\n    df = spark.read.schema(schema).option('header',True).option('mode','FAILFAST').csv(path)\n    \n    if limit is not None:\n        df = df.limit(limit)\n        \n    return df\n\ntrain_df = read_file('/kaggle/input/amex-default-prediction/train_data.csv')\ntest_df  = read_file('/kaggle/input/amex-default-prediction/test_data.csv')", "metadata": {"execution": {"iopub.status.busy": "2022-08-02T00:02:28.900168Z", "iopub.execute_input": "2022-08-02T00:02:28.900908Z", "iopub.status.idle": "2022-08-02T00:02:32.365531Z", "shell.execute_reply.started": "2022-08-02T00:02:28.900865Z", "shell.execute_reply": "2022-08-02T00:02:32.364086Z"}, "trusted": true}, "execution_count": null, "outputs": []}, {"cell_type": "code", "source": "#train_df.printSchema()", "metadata": {"execution": {"iopub.status.busy": "2022-08-02T00:02:32.367082Z", "iopub.execute_input": "2022-08-02T00:02:32.367595Z", "iopub.status.idle": "2022-08-02T00:02:32.381607Z", "shell.execute_reply.started": "2022-08-02T00:02:32.367543Z", "shell.execute_reply": "2022-08-02T00:02:32.379921Z"}, "trusted": true}, "execution_count": null, "outputs": []}, {"cell_type": "markdown", "metadata": {}, "source": "## 4. Memory Optimization: Downcasting\nMany categorical columns contain very small numbers (e.g., 1 to 5). Structurally storing these as 64-bit strings or floats is incredibly wasteful. We specifically target and downcast these low-cardinality features into strict `ByteType()`, saving enormous RAM overhead."}, {"cell_type": "markdown", "metadata": {}, "source": "## 4. Memory Optimization: Downcasting\nMany categorical columns contain very small numbers (e.g., 1 to 5). Structurally storing these as 64-bit strings or floats is incredibly wasteful. We specifically target and downcast these low-cardinality features into strict `ByteType()`, saving enormous RAM overhead."}, {"cell_type": "code", "source": "def cast_categ_2_byte(df):\n    for col_nm in cat_col-cat_lab_col:\n        df = df.withColumn(col_nm,df[col_nm].cast(pyspark.sql.types.ByteType()))\n    return df\n                           \ntrain_df = cast_categ_2_byte(train_df)\ntest_df = cast_categ_2_byte(test_df)", "metadata": {"execution": {"iopub.status.busy": "2022-08-02T00:02:32.387218Z", "iopub.execute_input": "2022-08-02T00:02:32.388195Z", "iopub.status.idle": "2022-08-02T00:02:33.999477Z", "shell.execute_reply.started": "2022-08-02T00:02:32.388153Z", "shell.execute_reply": "2022-08-02T00:02:33.998091Z"}, "trusted": true}, "execution_count": null, "outputs": []}, {"cell_type": "markdown", "metadata": {}, "source": "## 5. StringIndexer for Labels\nMachine learning models require numerical inputs. For the remaining string-based categorical columns (`D_63`, `D_64`), we use Spark ML's `StringIndexer` to strictly map ordinal/nominal text into optimized Byte indices before saving."}, {"cell_type": "markdown", "metadata": {}, "source": "## 5. StringIndexer for Labels\nMachine learning models require numerical inputs. For the remaining string-based categorical columns (`D_63`, `D_64`), we use Spark ML's `StringIndexer` to strictly map ordinal/nominal text into optimized Byte indices before saving."}, {"cell_type": "code", "source": "from pyspark.ml.feature import StringIndexerModel\n\nstring_indexer = StringIndexerModel.from_arrays_of_labels( \n    [['CL', 'CO', 'CR', 'XL', 'XM', 'XZ'], ['-1', 'O', 'R', 'U']],\n    inputCols=['D_63', 'D_64'],\n    outputCols=['D_63_index', 'D_64_index'],\n    handleInvalid=\"keep\")\n\ndef label_encoding(df):\n    df = string_indexer.transform(df)\n    for col_nm in cat_lab_col:\n        index_col = F.col(col_nm+'_index')\n        index_col = F.when(F.col(col_nm).isNull(),None).otherwise(index_col)\n        index_col = index_col.cast(pyspark.sql.types.ByteType())\n        df = df.withColumn(col_nm + '_index', index_col)\n        df = df.drop(col_nm)\n    return df\n\ntrain_df = label_encoding(train_df)\ntest_df = label_encoding(test_df)", "metadata": {"execution": {"iopub.status.busy": "2022-08-02T00:02:34.000993Z", "iopub.execute_input": "2022-08-02T00:02:34.001494Z", "iopub.status.idle": "2022-08-02T00:02:35.251096Z", "shell.execute_reply.started": "2022-08-02T00:02:34.001442Z", "shell.execute_reply": "2022-08-02T00:02:35.249657Z"}, "trusted": true}, "execution_count": null, "outputs": []}, {"cell_type": "code", "source": "#train_df.show(1,vertical=True)", "metadata": {"execution": {"iopub.status.busy": "2022-08-02T00:02:35.25294Z", "iopub.execute_input": "2022-08-02T00:02:35.25439Z", "iopub.status.idle": "2022-08-02T00:02:35.260298Z", "shell.execute_reply.started": "2022-08-02T00:02:35.254324Z", "shell.execute_reply": "2022-08-02T00:02:35.259001Z"}, "trusted": true}, "execution_count": null, "outputs": []}, {"cell_type": "markdown", "metadata": {}, "source": "## 6. Distributed Parquet Output\n**Why Parquet?** Parquet is a columnar storage format. It automatically compresses similar data heavily and allows models to ONLY read the columns they need (Predicate Pushdown).\n\n**Why Partition by `customer_ID`?** In downstream tasks (like feature engineering), we will need to aggregate time-series rows for each customer. By repartitioning the parquet files on `customer_ID` *now*, we ensure all records for a given customer live on the same worker node, completely eliminating network shuffling during modeling!"}, {"cell_type": "markdown", "metadata": {}, "source": "## 6. Distributed Parquet Output\n**Why Parquet?** Parquet is a columnar storage format. It automatically compresses similar data heavily and allows models to ONLY read the columns they need (Predicate Pushdown).\n\n**Why Partition by `customer_ID`?** In downstream tasks (like feature engineering), we will need to aggregate time-series rows for each customer. By repartitioning the parquet files on `customer_ID` *now*, we ensure all records for a given customer live on the same worker node, completely eliminating network shuffling during modeling!"}, {"cell_type": "code", "source": "def write_parquet(df,name,partition):\n    df.repartition(partition,'customer_ID')\\\n    .write\\\n    .parquet(\"amexparq/\" + name)\n\nwrite_parquet(train_df, \"train\", 20)\nwrite_parquet(test_df, \"test\", 20)", "metadata": {"execution": {"iopub.status.busy": "2022-08-02T00:02:35.262296Z", "iopub.execute_input": "2022-08-02T00:02:35.263357Z", "iopub.status.idle": "2022-08-02T00:20:23.659103Z", "shell.execute_reply.started": "2022-08-02T00:02:35.2633Z", "shell.execute_reply": "2022-08-02T00:20:23.644485Z"}, "trusted": true}, "execution_count": null, "outputs": []}]}