{"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":"## This notebook explains methods to prevent memory errors when creating a Covisitation Matrix\nI've seen a few people say they are getting memory error when creating a Covisitation Matrix.\nI am a beginner myself, and had to spend hours to understand how others did it. I tried to break down the techniques used from other notebooks as easy as possible, so if you have any questions or issues please leave a comment!\nThis notebook uses techniques and codes mainly from @cdeotte and @radek1 's notebook.\n\n- This notebook is in three steps:\n    - Technique 1: Read file in batches from parquet/csv ( 28 `df`s. It becomes 36 if you use [radek's full data](https://www.kaggle.com/datasets/radek1/otto-full-optimized-memory-footprint))\n    - Technique 2: Separate `aid_x` into ranges to calculate the score. (Read 28 `df`s 4 times each, with range of `aid_x` set differently)\n    - Technique 3: Combine many `df`s at a time for speedup.\n    \n    \n\n- After you read: You will understand many techniques to solve memory shortage.\n\n- Before you read: It's better that you know how to create a Covisitation Matrix. In this notebook, I simply use code snippets from @cdeotte's notebook with only minor tweaks in order to generate a Covisitation matrix.","metadata":{}},{"cell_type":"code","source":"# imports\nimport matplotlib.pyplot as plt\nimport pandas as pd, numpy as np\nfrom tqdm.notebook import tqdm\nimport os, sys, pickle, glob, gc\nfrom collections import Counter\nimport cudf, itertools\nimport pyarrow.parquet as pq\nprint('We will use RAPIDS version',cudf.__version__)","metadata":{"execution":{"iopub.status.busy":"2022-12-10T09:31:06.097818Z","iopub.execute_input":"2022-12-10T09:31:06.098194Z","iopub.status.idle":"2022-12-10T09:31:06.104739Z","shell.execute_reply.started":"2022-12-10T09:31:06.098161Z","shell.execute_reply":"2022-12-10T09:31:06.103680Z"},"trusted":true},"execution_count":null,"outputs":[]},{"cell_type":"markdown","source":"## Technique 1: First, let's try to read df in batches.\nThere are actually two ways to read file in batches (same as chunks, but I will say batches only for datafframes throughout the document for less confusion).\n\n1. Using a dataset that is already divided into parts. [This dataset](https://www.kaggle.com/datasets/columbia2131/otto-chunk-data-inparquet-format) is an example.\n2. Using built-in functions in pandas or pyarrow to read a single file in batches.\n\n[Rerank Model by Deotte](https://www.kaggle.com/code/cdeotte/candidate-rerank-model-lb-0-575/notebook) uses multiple parquet files in the first place, which means\nthat they are already in batches. However, in this notebook, we will use dataset made\nby @radek, which makes a code a bit more concise. We need to use `pyarrow.parquet` in order to  easily\nread parquet files in chunks.","metadata":{}},{"cell_type":"code","source":"PATH = '../input/otto-train-and-test-data-for-local-validation/'\n# pyarrow.parquet IS IMPORTED AS pq.\ntrain_parquet_file = pq.ParquetFile(PATH + 'train.parquet')\ntest_parquet_file = pq.ParquetFile(PATH + 'test.parquet')\ndf_batches = []\n# YOU CAN TUNE batch_size TO YOUR NEEDS. IT REPRESENTS THE ROW NUMBER FOR EACH BATCH.\nfor batch in train_parquet_file.iter_batches(batch_size=65536*100):\n    df_batches.append(batch.to_pandas())\nfor batch in test_parquet_file.iter_batches(batch_size=65536*100):\n    df_batches.append(batch.to_pandas())","metadata":{"execution":{"iopub.status.busy":"2022-12-10T09:31:07.935286Z","iopub.execute_input":"2022-12-10T09:31:07.936257Z","iopub.status.idle":"2022-12-10T09:31:21.393446Z","shell.execute_reply.started":"2022-12-10T09:31:07.936217Z","shell.execute_reply":"2022-12-10T09:31:21.392446Z"},"trusted":true},"execution_count":null,"outputs":[]},{"cell_type":"code","source":"len(df_batches) # NOW WE HAVE DATASET DIVIDED INTO MULTIPLE PARTS","metadata":{"execution":{"iopub.status.busy":"2022-12-10T09:31:45.965936Z","iopub.execute_input":"2022-12-10T09:31:45.966608Z","iopub.status.idle":"2022-12-10T09:31:45.974656Z","shell.execute_reply.started":"2022-12-10T09:31:45.966570Z","shell.execute_reply":"2022-12-10T09:31:45.973519Z"},"trusted":true},"execution_count":null,"outputs":[]},{"cell_type":"markdown","source":"> _\"Now that we have separated data into 28 parts, maybe we can generate a Covistation Matrix!\"_\n\nSadly, but no. If you comment out `%%script echo skipping` part and run the below code, it causes a memory error. We need to split the data into **even smaller**\nparts.","metadata":{}},{"cell_type":"code","source":"%%script echo skipping\n\n# THIS CODE GENERATES A CLICK2CLICK COVISITATION MATRIX.\n# RUNNING THIS WILL GIVE A MEMORY ERROR.\nfor i, df_batch in enumerate(df_batches):\n    print(str(i) + '/' + str(len(df_batches)), end=' ')\n    df = cudf.DataFrame(df_batch)\n    df = df.loc[d['type'].isin([0])] # pickout clicks only\n    df = df.sort_values(['session', 'ts'], ascending=[True, False])\n    # TAIL\n    df = df.reset_index(drop=True)\n    df['n'] = df.groupby('session').cumcount()\n    df = df.loc[d.n<30].drop('n', axis=1)\n    # SELF JOIN\n    df = df.merge(d, on='session')\n    df = df.loc[((df.ts_x - df.ts_y).abs()< 24 * 60 * 60) & (d.aid_x != d.aid_y)]\n    # WEIGHTS\n    df = df[['session', 'aid_x', 'aid_y','ts_x']].drop_duplicates(['session', 'aid_x', 'aid_y'])\n    df['wgt'] = 1 + 3*(d.ts_x - 1659304800)/(1662328791-1659304800)\n    df = df[['aid_x','aid_y','wgt']]\n    df.wgt = df.wgt.astype('float32')\n    df = df.groupby(['aid_x','aid_y']).wgt.sum()\n    # ADD\n    if i == 0:\n        df_result = df\n    else:\n        df_result = df_result.add(df, fill_value=0)\n    ","metadata":{"execution":{"iopub.status.busy":"2022-12-10T09:32:25.426461Z","iopub.execute_input":"2022-12-10T09:32:25.426834Z","iopub.status.idle":"2022-12-10T09:32:25.501327Z","shell.execute_reply.started":"2022-12-10T09:32:25.426804Z","shell.execute_reply":"2022-12-10T09:32:25.500224Z"},"trusted":true},"execution_count":null,"outputs":[]},{"cell_type":"markdown","source":"## Technique 2: Now, we use separate data into SMALLER CHUNKS with `aid_x`'s index\nI don't know who originally used this technique in the competition, but I copied this idea from @cdeotte's notebooks. Since we are generating the score of pairs of `aid_x` and `aid_y`, we can separate data with `aid_x`. For instance, let's say we are splitting index of `aid_x` into four parts. We know there are `1855602` number of aids, so both `aid_x` and `aid_y` has range of `0~1855602`. Separating into four parts would be `0~465000`, `465000~930000`, `930000~1395000`, and `1395000~18655602`. For each one of the four parts, we calculate the score for pairs of `aid_x` and `aid_y`. Then, we add them together by `df1.add(df2)` multiple times. This adds the scores for unique `aid_x` and `aid_y` pair, using much less memory at a time.\n\nThe code that separates `aid_x` is `df = df.loc[(df.aid_x >= PART*SIZE)&(df.aid_x < (PART+1)*SIZE)]`, which is used in the below cells. This allows us to later add the scores, but save significant portion of the memory space used at once.","metadata":{}},{"cell_type":"code","source":"aid_max = 0\nfor df in df_batches:\n    aid_max = max(aid_max, np.max(df.aid))\nprint(aid_max)","metadata":{"execution":{"iopub.status.busy":"2022-12-10T09:45:44.473585Z","iopub.execute_input":"2022-12-10T09:45:44.473939Z","iopub.status.idle":"2022-12-10T09:45:44.613727Z","shell.execute_reply.started":"2022-12-10T09:45:44.473906Z","shell.execute_reply":"2022-12-10T09:45:44.612666Z"},"trusted":true},"execution_count":null,"outputs":[]},{"cell_type":"code","source":"# THE MAXIMUM `aid` NUMBER IS LESS THAN 1.86e6.\nDISK_PIECES = 4\nSIZE = 1.86e6/DISK_PIECES","metadata":{"execution":{"iopub.status.busy":"2022-12-10T09:45:44.705913Z","iopub.execute_input":"2022-12-10T09:45:44.706210Z","iopub.status.idle":"2022-12-10T09:45:44.712160Z","shell.execute_reply.started":"2022-12-10T09:45:44.706183Z","shell.execute_reply":"2022-12-10T09:45:44.710101Z"},"trusted":true},"execution_count":null,"outputs":[]},{"cell_type":"code","source":"%%time\n\n# GENERATE A CLICK2CLICK COVISITATION MATRIX\nfor PART in range(DISK_PIECES):\n    print(f'PART{PART}')\n    for i, df_batch in enumerate(df_batches):\n        print(f'{i}/{len(df_batches)}', end=' ')\n        df = cudf.DataFrame(df_batch)\n        df = df.loc[df['type'].isin([0])] # pickout clicks only\n        df = df.sort_values(['session', 'ts'], ascending=[True, False])\n        # TAIL\n        df = df.reset_index(drop=True)\n        df['n'] = df.groupby('session').cumcount()\n        df = df.loc[df.n<30].drop('n', axis=1)\n        # SELF JOIN\n        df = df.merge(df, on='session')\n        df = df.loc[((df.ts_x - df.ts_y).abs()< 24 * 60 * 60) & (df.aid_x != df.aid_y)]\n        \n        #### USE DISK_PIECES ####\n        # USE aid_x AS AN INDEX TO SPERATE DATA INTO EVEN SMALLER CHUNKS\n        df = df.loc[(df.aid_x >= PART*SIZE)&(df.aid_x < (PART+1)*SIZE)]\n        \n        # WEIGHTS\n        df = df[['session', 'aid_x', 'aid_y','ts_x']].drop_duplicates(['session', 'aid_x', 'aid_y'])\n        df['wgt'] = 1 + 3*(df.ts_x - 1659304800)/(1662328791-1659304800)\n        df = df[['aid_x','aid_y','wgt']]\n        df.wgt = df.wgt.astype('float32')\n        df = df.groupby(['aid_x','aid_y']).wgt.sum()\n        # ADD\n        if i == 0:\n            df_result = df\n        else:\n            df_result = df_result.add(df, fill_value=0)\n    # DON'T MERGE WITH OTHER DISK_PIECES, BUT DIRECTLY SAVE THE RESULT.\n    # YOU CANNOT MERGE THEM DUE TO RAM SHORTAGE.\n    # SO WHAT YOU DO IS, YOU LATER MERGE THE SEPARATE FILES AND THEN GENERATE CANDIDATES.\n    \n    # POST PROCESSING BEFORE SAVING\n    df_result = df_result.reset_index()\n    df_result = df_result.sort_values(['aid_x','wgt'],ascending=[True,False])\n    df_result = df_result.reset_index(drop=True)\n    df_result['n'] = df_result.groupby('aid_x').aid_y.cumcount()\n    df_result = df_result.loc[df_result.n<30].drop('n',axis=1) # SAVE TOP 30\n    \n    # SAVE PART TO DISK (convert to pandas first uses less memory)\n    df_result.to_pandas().to_parquet(f'TOP30_CL2CL_{PART}.pqt')\n    \n    del df_result\n    gc.collect()\n    print()","metadata":{"execution":{"iopub.status.busy":"2022-12-10T10:52:04.910537Z","iopub.execute_input":"2022-12-10T10:52:04.910880Z","iopub.status.idle":"2022-12-10T10:55:27.023854Z","shell.execute_reply.started":"2022-12-10T10:52:04.910848Z","shell.execute_reply":"2022-12-10T10:55:27.022772Z"},"trusted":true},"execution_count":null,"outputs":[]},{"cell_type":"markdown","source":"**Technique 1 and 2 is enough to get through the memory requirement! Technique 3 allows you to boost up the add speed.**","metadata":{}},{"cell_type":"code","source":"# REMOVE PREVIOUSLY GENERATED DATA\nfor i in range(0, 4):\n    os.remove(f\"/kaggle/working/TOP30_CL2CL_{i}.pqt\")","metadata":{},"execution_count":null,"outputs":[]},{"cell_type":"markdown","source":"## Technique 3: Process many `df`s at a time \nBy processing many `df`s at a time, You spend less computational resource on `add` operation between `df`s. ","metadata":{}},{"cell_type":"code","source":"%%time\n# 118\nfor PART in range(DISK_PIECES):\n    print(f'PART{PART}')\n    \n    for k in range(28//6+1): # same as `range(len(df_batches//6) + 1)`\n        start = k*6\n        end = (k+1)*6\n        for i, df_batch in enumerate(df_batches[start:end]):\n            print(f'{k*6 + i + 1}/{len(df_batches)}', end=' ')\n            df = cudf.DataFrame(df_batch)\n            df = df.loc[df['type'].isin([0])] # pickout clicks only\n            df = df.sort_values(['session', 'ts'], ascending=[True, False])\n            # TAIL\n            df = df.reset_index(drop=True)\n            df['n'] = df.groupby('session').cumcount()\n            df = df.loc[df.n<30].drop('n', axis=1)\n            # SELF JOIN\n            df = df.merge(df, on='session')\n            df = df.loc[((df.ts_x - df.ts_y).abs()< 24 * 60 * 60) & (df.aid_x != df.aid_y)]\n\n            #### USE DISK_PIECES ####\n            # USE aid_x AS AN INDEX TO SPERATE DATA INTO EVEN SMALLER CHUNKS\n            df = df.loc[(df.aid_x >= PART*SIZE)&(df.aid_x < (PART+1)*SIZE)]\n\n            # WEIGHTS\n            df = df[['session', 'aid_x', 'aid_y','ts_x']].drop_duplicates(['session', 'aid_x', 'aid_y'])\n            df['wgt'] = 1 + 3*(df.ts_x - 1659304800)/(1662328791-1659304800)\n            df = df[['aid_x','aid_y','wgt']]\n            df.wgt = df.wgt.astype('float32')\n            df = df.groupby(['aid_x','aid_y']).wgt.sum()\n            \n            # ADD\n            if i == 0:\n                df_inner = df\n            else:\n                df_inner = df_inner.add(df, fill_value=0)\n        if k == 0:\n            df_result = df_inner\n        else:\n            df_result = df_result.add(df_inner, fill_value=0)\n        \n        del df_inner\n        gc.collect()\n            \n    # POST PROCESSING BEFORE SAVING\n    df_result = df_result.reset_index()\n    df_result = df_result.sort_values(['aid_x','wgt'],ascending=[True,False])\n    df_result = df_result.reset_index(drop=True)\n    df_result['n'] = df_result.groupby('aid_x').aid_y.cumcount()\n    df_result = df_result.loc[df_result.n<30].drop('n',axis=1) # SAVE TOP 30\n    \n    # SAVE PART TO DISK (convert to pandas first uses less memory)\n    df_result.to_pandas().to_parquet(f'TOP30_CL2CL_{PART}.pqt')\n    \n    del df_result\n    gc.collect()\n    print()","metadata":{"execution":{"iopub.status.busy":"2022-12-10T10:49:19.911863Z","iopub.execute_input":"2022-12-10T10:49:19.912243Z","iopub.status.idle":"2022-12-10T10:51:38.955759Z","shell.execute_reply.started":"2022-12-10T10:49:19.912209Z","shell.execute_reply":"2022-12-10T10:51:38.954772Z"},"trusted":true},"execution_count":null,"outputs":[]},{"cell_type":"code","source":"","metadata":{},"execution_count":null,"outputs":[]}]}