{"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":"[AMEX Default Prediction Competition](https://www.kaggle.com/competitions/amex-default-prediction) - dask pre-processing\n===================================\n**Author:** [@xaviernogueira](https://github.com/xaviernogueira)\n\n**Purpose:** Leverage [`dask`](https://docs.dask.org/en/stable/) to re-format [raw AMEX train/test .csv data](https://www.kaggle.com/competitions/amex-default-prediction/data?select=train_data.csv) column dtypes and export to .parquet files in a locally distribted manner, avoiding poor pandas performance and RAM errors.\n\n**Method notes:**\n* We use [`dask.distributed`](https://distributed.dask.org/en/stable/) on out local machine to parallelize computation across 4 workers.\n* We convert all object columns (including `customer_ID`), as well as other hardcoded columns, into the category dtype to optimize perfomance.\n* We convert column `S_2`, which stores statement dates, into the `datetime64[ns]` dtype.\n* As opposed to pandas, in dask `df.set_index()` is [very computationally expensive](https://coiled.io/blog/dask-set-index-dataframe/), and therefore is avoided for now.\n* The output data is partitioned and saved into many smaller .parquet files. These can easily be bacth read into a single dask dataframe via [dask.dataframe.read_parquet(DIR_NAME)](https://docs.dask.org/en/latest/generated/dask.dataframe.read_parquet.html). Alternatively, individual partitians could be used to build out workflows rapidly.\n* This workflow took approximately 12 minutes on a 4 core, 8 thread 11th Gen Intel(R) Core(TM) i7-1185G7 processer with 24GB RAM.\n\n![DASK](https://docs.dask.org/en/latest/_images/dask_horizontal.svg)","metadata":{"tags":[]}},{"cell_type":"markdown","source":"## Import dependencies","metadata":{"jp-MarkdownHeadingCollapsed":true,"tags":[]}},{"cell_type":"code","source":"# import base os dependencies\nfrom pathlib import Path\nimport os\nimport datetime\nimport gc\n\n# import performance dependencies (dask and numba)\nimport dask.dataframe as dd\nfrom dask.distributed import Client\nimport pyarrow","metadata":{},"execution_count":null,"outputs":[]},{"cell_type":"markdown","source":"## Start `dask.distributed` client","metadata":{"jp-MarkdownHeadingCollapsed":true,"tags":[]}},{"cell_type":"code","source":"# get the starttime so we can assess notebook perfomance\nstart_time = datetime.datetime.now()\nprint(f'Start time: {start_time}')\n\n# get the dask client working (local machine), we see this gives us access to 4 workers!\nclient = Client()\nclient","metadata":{},"execution_count":null,"outputs":[]},{"cell_type":"markdown","source":"# Use dask to import and pre-process data\nHere we do the following:\n1. Import the test_data, train_data, and train_labels .csv files as `dask` dataframes.\n2. Convert the `S_2` feature to from object dtype to `datetime64[ns]`.\n3. Format object columns into categorical dtype for performance (except `customer_ID`).\n4. Output to partitianed .parquet files using `pyarrow` and `gzip` compressions. The idea here is that we can commit the .parquet data files to GitHub without using [Git LFS](https://git-lfs.github.com/) and preserve the updated column dtypes.","metadata":{"tags":[]}},{"cell_type":"markdown","source":"## Imported .csv data to dask","metadata":{"jp-MarkdownHeadingCollapsed":true,"tags":[]}},{"cell_type":"markdown","source":"### Find all .csv files in DATA_FOLDER and import to dataframes","metadata":{"tags":[]}},{"cell_type":"code","source":"# Confirm your current working directory (cwd)\nprint(f'Working directory: {Path.cwd()}')\n\n# SET YOUR DATA FOLDER (where .csv or .parquet AMEX data files are stores)\nDATA_FOLDER = str(Path.cwd()) + '\\\\data_by_competition\\\\amex-default-prediction'\nprint(f'Data directory: {DATA_FOLDER}')\n\n# define this here\ngroup_by_cols = ['customer_ID', 'S_2']","metadata":{},"execution_count":null,"outputs":[]},{"cell_type":"code","source":"def id_csvs(data_dir: str):\n    \"\"\"\"Takes a dir path containing .csv tables, and returns a dictionary with the .csv names as keys, with empty lists as items\"\"\"\n    csv_dict = {}\n\n    # iterate over data_dir looking for .csv files\n    for file in os.listdir(data_dir):\n        if file[-4:] == '.csv':\n            # add to csv_dict\n            csv_dict[file.split('.', 2)[0]] = f'{data_dir}\\\\{file}'\n\n    # print csv_dict items, and return the dictionary\n    print(f'csv_dict initiated {csv_dict.keys()}')\n    return csv_dict\n\n\n# define a pre-compile a function to import data from csv\ndef import_csv(data_dir: str, csv_dict: dict, index_col: str = None):\n    \"\"\"\n    Uses dask to partitian into multiple dask.dataframes with max # of rows = chunk_rise\n    :param data_dir:\n    :param csv_dict: the output of id_csv\n    :param index_col: (optional) index column for the csv\n    :returns: csv_dict, but updated so that each key stores a the dask dataframe\n    \"\"\"\n    for data_name in csv_dict.keys():\n        print(f'Importing {data_name}...')\n        csv_path = csv_dict[data_name]\n\n        # read the csv sampling 5MB at a time to determin dtypes\n        df = dd.read_csv(csv_path, engine='c', sample=5000000)\n\n        # if an index is provided, warn the user about the performance penalty\n        if index_col is not None:\n            print(f'WARNING: Setting index to {index_col} will drastically hurt performance!')\n            df = df.set_index(index_col)\n\n        # update the dictionary and collect trash\n        sub_dict = {data_name: df}\n        csv_dict.update(sub_dict)\n        gc.collect()\n    return csv_dict","metadata":{},"execution_count":null,"outputs":[]},{"cell_type":"code","source":"%%time\n# make csv_dict storing ex: {'test_data': 'TEST//DATA//PATH.csv', ...}\ncsv_dict = id_csvs(DATA_FOLDER)\n\n# update csv_dict to store dask dataframes for each data type\nimport_csv(csv_dict=csv_dict, data_dir=DATA_FOLDER)\n_ = gc.collect()\nprint(f'csv_dict stores dask.dataframes with the following keys: {csv_dict.keys()}')","metadata":{},"execution_count":null,"outputs":[]},{"cell_type":"markdown","source":"## Convert dtypes ","metadata":{"tags":[]}},{"cell_type":"markdown","source":"### Convert object columns to category dtype","metadata":{"tags":[]}},{"cell_type":"code","source":"# make a function that convert object columns to category dtype\ndef object_to_cat(df, ignore_cols: list = []):\n    \"\"\" Converts all object columns to category type in dask.dataframe, unless specified in ignore_cols\"\"\"\n    object_cols = df.select_dtypes(include='object')\n\n    # for each object column, convert the original dataframes dtype\n    for col in object_cols:\n        if str(col) not in ignore_cols:\n            df[col] = df[col].astype('category')\n\n    return df","metadata":{},"execution_count":null,"outputs":[]},{"cell_type":"code","source":"%%time\n# apply the function to the data dataframes (no train_labels)\nfor title in csv_dict.keys():\n    if 'data' in title:\n        df = csv_dict[title]\n        \n        # convert all other object columns to category\n        df = object_to_cat(df, ignore_cols=['S_2'])","metadata":{},"execution_count":null,"outputs":[]},{"cell_type":"markdown","source":"### Convert column `S_2` to `datetime64[ns]` dtype","metadata":{}},{"cell_type":"code","source":"%%time\n# apply the functions to the data dataframes (no train_labels)\nfor title in csv_dict.keys():\n    if 'data' in title:\n        df = csv_dict[title]\n\n        # convert S_2 to datetime\n        df['S_2'] = dd.to_datetime(df['S_2'], format='%Y-%m-%d')","metadata":{},"execution_count":null,"outputs":[]},{"cell_type":"code","source":"_ = gc.collect()\ndf.info()","metadata":{},"execution_count":null,"outputs":[]},{"cell_type":"markdown","source":"### Convert columns to category dtype (using hardcoded columns headers)","metadata":{"tags":[]}},{"cell_type":"code","source":"# hardcode all categorical columns\nhardcoded_cat_cols = [\"B_30\", \"B_38\", \"D_114\", \"D_116\", \"D_117\", \"D_120\", \"D_126\", \"D_63\", \"D_64\", \"D_66\", \"D_68\"]\n\n# get a list of all columns\nall_cols = [c for c in list(csv_dict['train_data'].columns) if c not in group_by_cols and c in list(csv_dict['test_data'].columns)]","metadata":{},"execution_count":null,"outputs":[]},{"cell_type":"code","source":"%%time\n# convert any unconverted categorical columns to category dtype for performance\nfor title in csv_dict.keys():\n    if 'data' in title:\n        df = csv_dict[title]\n        df[hardcoded_cat_cols] = df[hardcoded_cat_cols].astype('category')","metadata":{},"execution_count":null,"outputs":[]},{"cell_type":"code","source":"df.info()","metadata":{},"execution_count":null,"outputs":[]},{"cell_type":"markdown","source":"### Get dtype column lists (`num_cols` and `cat_cols`)","metadata":{}},{"cell_type":"code","source":"# get an updated list of catgeorical and numverical columns columns\ncat_cols = list(csv_dict['train_data'].select_dtypes('category').columns)\ncat_cols.sort()\nprint(f'category dtype columns (n={len(cat_cols)}): {cat_cols}')\n\nnum_cols = [i for i in all_cols if i not in cat_cols]\nnum_cols.sort()\nprint(f'numerical dtype columns (n={len(num_cols)}): {num_cols}')","metadata":{},"execution_count":null,"outputs":[]},{"cell_type":"code","source":"# collect garbage\n_ = gc.collect()","metadata":{},"execution_count":null,"outputs":[]},{"cell_type":"code","source":"# see if the output looks reasonable\ncsv_dict['train_data'].info()\ncsv_dict['train_data'].head(n=3)","metadata":{},"execution_count":null,"outputs":[]},{"cell_type":"markdown","source":"# Save dataframes to .parquet\nHere a new folder is made within `DATA_FOLDER` which stores .parquet partitial files for each AMEX csv. For example, `DATA_FOLDER/train_data/` will contain many files of the fort `part.212.parquet`. Finally we calculate our notebook's run time for comparison against pandas workflows.","metadata":{"tags":[]}},{"cell_type":"code","source":"%%time\n# convert all dataframes to parquets\nfor title in csv_dict.keys():\n    out_path = DATA_FOLDER + f'\\\\{title}_parquets'\n    df = csv_dict[title]\n    \n    # convert the dask dataframe to parquet\n    dd.to_parquet(df, out_path, engine='pyarrow', compression='gzip')\n    print(f'{title} saved @ {out_path}')\n    \n    # collect garbage and delete the previous dask dataframe\n    _ = gc.collect()\n    del df\n    \nprint(f'Done! all .parquet files saved @ {DATA_FOLDER}')","metadata":{},"execution_count":null,"outputs":[]},{"cell_type":"code","source":"# assess notebook run time\nend_time = datetime.datetime.now()\nrun_time = (end_time - start_time).total_seconds()\nprint(f'End time: {start_time}\\n')\nprint(f'******* NOTEBOOK RUNTIME: {run_time} seconds *******')","metadata":{},"execution_count":null,"outputs":[]}]}