{"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":"# Version details\n- V2: fulltype: weight 163 | click: time weight | purchase: 7days (apply duplicate removal for all)\n- V3: fulltype: weight 163 | click: time weight | purchase: 14days (apply duplicate removal for all)\n- V4: 4 covisitation","metadata":{}},{"cell_type":"code","source":"# import pandas as pd\n\n# train_df = pd.read_parquet(\"/kaggle/input/otto-train-and-test-data-for-local-validation/train.parquet\")\n# test_df = pd.read_parquet(\"/kaggle/input/otto-train-and-test-data-for-local-validation/test.parquet\")\n# label_df = pd.read_parquet(\"/kaggle/input/otto-train-and-test-data-for-local-validation/test_labels.parquet\")\n\n# print(f\"Len: {len(train_df)} - {len(test_df)} - {len(label_df)}\")\n# print(f\"Min session: {train_df.session.min()} - {test_df.session.min()} - {label_df.session.min()}\")\n# print(f\"Max session: {train_df.session.max()} - {test_df.session.max()} - {label_df.session.max()}\")\n# print(f\"Num session: {train_df.session.nunique()} - {test_df.session.nunique()} - {label_df.session.nunique()}\")","metadata":{"execution":{"iopub.status.busy":"2022-12-11T07:24:34.577561Z","iopub.execute_input":"2022-12-11T07:24:34.578027Z","iopub.status.idle":"2022-12-11T07:24:34.598078Z","shell.execute_reply.started":"2022-12-11T07:24:34.577927Z","shell.execute_reply":"2022-12-11T07:24:34.597224Z"},"trusted":true},"execution_count":null,"outputs":[]},{"cell_type":"code","source":"# VER = 5\nDEBUG = False\nTYPE_WEIGHT_FULLTYPE = {0:1, 1:6, 2:3}\nTYPE_WEIGHT_PURCHASE = {0:1, 1:1, 2:1}\nTYPE_WEIGHT_CLICK = None # use time weight\nDELTA_TS_FULLTYPE = 24\nDELTA_TS_CLICK = 24 # hours\nDELTA_TS_BUY = 24*14\nLENGTH_TAIL_SESSION = 30\nTOP_K = 20\nSAVE_DICT_FORMAT = False\nLEAK_DATA = True\n\nMIN_CLICK = 1.0\nMIN_CO = 1.0\nMIN_TS = 1659304800\nMAX_TS = 1661723996 if LEAK_DATA else 1661119199","metadata":{"execution":{"iopub.status.busy":"2022-12-11T07:24:34.599916Z","iopub.execute_input":"2022-12-11T07:24:34.600362Z","iopub.status.idle":"2022-12-11T07:24:34.611577Z","shell.execute_reply.started":"2022-12-11T07:24:34.600325Z","shell.execute_reply":"2022-12-11T07:24:34.610380Z"},"trusted":true},"execution_count":null,"outputs":[]},{"cell_type":"code","source":"import pandas as pd, numpy as np\nfrom tqdm.notebook import tqdm\nimport os, sys, pickle, glob, gc\nfrom collections import Counter\nimport cudf, itertools\nprint('We will use RAPIDS version',cudf.__version__)","metadata":{"execution":{"iopub.status.busy":"2022-12-11T07:24:34.613631Z","iopub.execute_input":"2022-12-11T07:24:34.613987Z","iopub.status.idle":"2022-12-11T07:24:37.814693Z","shell.execute_reply.started":"2022-12-11T07:24:34.613954Z","shell.execute_reply":"2022-12-11T07:24:37.813229Z"},"trusted":true},"execution_count":null,"outputs":[]},{"cell_type":"code","source":"%%time\n\n# CACHE FUNCTIONS\ndef read_file(f):\n    return cudf.DataFrame( data_cache[f] )\n\ndef read_file_to_cache(f):\n    df = pd.read_parquet(f)\n    df.ts = (df.ts/1000).astype('int32')\n    df['type'] = df['type'].map(type_labels).astype('int8')\n    return df\n\n# CACHE THE DATA ON CPU BEFORE PROCESSING ON GPU\ndata_cache = {}\ntype_labels = {'clicks':0, 'carts':1, 'orders':2}\nif LEAK_DATA:\n    files = glob.glob('/kaggle/input/otto-validation/*_parquet/*')\nelse:\n    files = glob.glob('/kaggle/input/otto-validation/train_parquet/*')\nfor f in tqdm(files): \n    data_cache[f] = read_file_to_cache(f)\n\n# CHUNK PARAMETERS\nREAD_CT = 5 # sub-chunk\nCHUNK = int( np.ceil( len(files)/6 ))\nprint(f'We will process {len(files)} files, in groups of {READ_CT} and chunks of {CHUNK}.')","metadata":{"execution":{"iopub.status.busy":"2022-12-11T07:29:42.683127Z","iopub.execute_input":"2022-12-11T07:29:42.683807Z","iopub.status.idle":"2022-12-11T07:30:28.052863Z","shell.execute_reply.started":"2022-12-11T07:29:42.683769Z","shell.execute_reply":"2022-12-11T07:30:28.051750Z"},"trusted":true},"execution_count":null,"outputs":[]},{"cell_type":"code","source":"suffixes = ['B', 'KB', 'MB', 'GB', 'TB', 'PB']\ndef humansize(nbytes):\n    i = 0\n    while nbytes >= 1024 and i < len(suffixes)-1:\n        nbytes /= 1024.\n        i += 1\n    f = ('%.2f' % nbytes).rstrip('0').rstrip('.')\n    return '%s %s' % (f, suffixes[i])","metadata":{"execution":{"iopub.status.busy":"2022-12-11T07:30:28.055019Z","iopub.execute_input":"2022-12-11T07:30:28.055706Z","iopub.status.idle":"2022-12-11T07:30:28.063579Z","shell.execute_reply.started":"2022-12-11T07:30:28.055667Z","shell.execute_reply":"2022-12-11T07:30:28.062338Z"},"trusted":true},"execution_count":null,"outputs":[]},{"cell_type":"code","source":"# dfs = []\n# for file in files:\n#     dfs.append(data_cache[file])\n# dfs = pd.concat(dfs, ignore_index=True)\n# dfs[\"session\"] = dfs[\"session\"].astype(\"int32\")\n# dfs[\"aid\"] = dfs[\"aid\"].astype(\"int32\")","metadata":{"execution":{"iopub.status.busy":"2022-12-11T07:24:54.693045Z","iopub.execute_input":"2022-12-11T07:24:54.693648Z","iopub.status.idle":"2022-12-11T07:24:54.709566Z","shell.execute_reply.started":"2022-12-11T07:24:54.693612Z","shell.execute_reply":"2022-12-11T07:24:54.708153Z"},"trusted":true},"execution_count":null,"outputs":[]},{"cell_type":"code","source":"# humansize(dfs.memory_usage(index=True, deep=True).sum())","metadata":{"execution":{"iopub.status.busy":"2022-12-11T07:24:54.711353Z","iopub.execute_input":"2022-12-11T07:24:54.712155Z","iopub.status.idle":"2022-12-11T07:24:54.719957Z","shell.execute_reply.started":"2022-12-11T07:24:54.712108Z","shell.execute_reply":"2022-12-11T07:24:54.718923Z"},"trusted":true},"execution_count":null,"outputs":[]},{"cell_type":"markdown","source":"# Co-visitation: co-occurence\n- duration < 1 day\n- top 20/40\n- the last 30/? action per session","metadata":{}},{"cell_type":"code","source":"%%time\n# USE SMALLEST DISK_PIECES POSSIBLE WITHOUT MEMORY ERROR\nTYPE_COO = \"fulltype\"\nOUTPUT_NAME = f\"top_{TOP_K}_{TYPE_COO}_{DELTA_TS_FULLTYPE}hours\"\nDISK_PIECES = 4 \nSIZE = 1.86e6/DISK_PIECES # total session\n\n# COMPUTE IN PARTS FOR MEMORY MANGEMENT\n\n# 146 files -> 6 CHUNK SIZE 25 -> 1 CHUNK: 5 GROUP SIZE 5\n# 1.86e6 aids -> 4 PART\n\nfor PART in range(DISK_PIECES):\n    print()\n    print('### DISK PART',PART+1)\n    \n    # MERGE IS FASTEST PROCESSING CHUNKS WITHIN CHUNKS\n    # => OUTER CHUNKS\n    for j in range(6):\n        a = j*CHUNK\n        b = min( (j+1)*CHUNK, len(files) )\n        print(f'\\nProcessing files {a} thru {b-1} in groups of {READ_CT}...', end=' ')\n        \n        # => INNER CHUNKS\n        for k in range(a,b,READ_CT):\n            # READ FILE\n            df = [read_file(files[k])]\n            for i in range(1,READ_CT): \n                if k+i<b: df.append( read_file(files[k+i]) )\n            df = cudf.concat(df,ignore_index=True,axis=0)\n            df = df.sort_values(['session','ts'],ascending=[True,False])\n            \n            # USE TAIL OF SESSION\n            df = df.reset_index(drop=True)\n            df['n'] = df.groupby('session').cumcount()\n            df = df.loc[df.n<LENGTH_TAIL_SESSION].drop('n',axis=1)\n            \n            # CREATE PAIRS\n            df = df.merge(df,on='session')\n            \n            df = df.loc[ ((df.ts_x - df.ts_y).abs()< DELTA_TS_FULLTYPE * 60 * 60) & (df.aid_x != df.aid_y) ]\n            \n            # MEMORY MANAGEMENT COMPUTE IN PARTS\n            df = df.loc[(df.aid_x >= PART*SIZE)&(df.aid_x < (PART+1)*SIZE)]\n            \n            # ASSIGN WEIGHTS\n            df = df.drop_duplicates(['session', 'aid_x', 'aid_y'])\n            df['wgt'] = df.type_y.map(TYPE_WEIGHT_FULLTYPE)\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            # COMBINE INNER CHUNKS\n            if k==a: tmp2 = df\n            else: tmp2 = tmp2.add(df, fill_value=0)\n            print(k,', ',end='')\n        \n        # COMBINE OUTER CHUNKS\n        if a==0: tmp = tmp2\n        else: tmp = tmp.add(tmp2, fill_value=0)\n        del tmp2, df\n        gc.collect()\n        \n        if DEBUG and j >= 1:\n            break\n    # CONVERT MATRIX TO DICTIONARY\n    tmp = tmp.reset_index()\n    tmp = tmp.sort_values(['aid_x','wgt'],ascending=[True,False])\n    \n    # SAVE TOP K\n    tmp = tmp.reset_index(drop=True)\n    tmp['n'] = tmp.groupby('aid_x').aid_y.cumcount()\n    tmp = tmp.loc[tmp.n<TOP_K].drop('n',axis=1)\n    tmp = tmp.reset_index(drop=True) # for saving memory\n    \n    # SAVE PART TO DISK (convert to pandas first uses less memory)\n    print(\"Saving sub_parquet...\", end ='')\n    tmp.to_pandas().to_parquet(f'/kaggle/working/{OUTPUT_NAME}_{PART}.pqt')\n    print(\"Saved!\")\n    del tmp\n    gc.collect()\n    if DEBUG and PART >= 1:\n        break\n\nprint(\"Saving full_parquet and full_dict...\", end='')\nfull_parquet = []\nfull_dict = dict()\nfor part in range(DISK_PIECES):\n    sub_parquet = pd.read_parquet(f'/kaggle/working/{OUTPUT_NAME}_{part}.pqt')\n    full_parquet.append(sub_parquet)\n    if SAVE_DICT_FORMAT:\n        full_dict.update(sub_parquet.groupby(\"aid_x\").apply(lambda df: Counter(dict(zip(df.aid_y, df.wgt)))))\n    os.remove(f'/kaggle/working/{OUTPUT_NAME}_{part}.pqt')\n    if DEBUG:\n        break\n    del sub_parquet\n    gc.collect()\nfull_parquet = pd.concat(full_parquet, axis=0, ignore_index=True)\nfull_parquet.to_parquet(f'/kaggle/working/{OUTPUT_NAME}_full.pqt')\nif SAVE_DICT_FORMAT:\n    with open(f'/kaggle/working/{OUTPUT_NAME}_full.pickle', 'wb') as file:\n        pickle.dump(full_dict, file)\nprint(\"Saved!\")\n\ndel full_parquet, full_dict\ngc.collect()","metadata":{"execution":{"iopub.status.busy":"2022-12-11T07:30:33.087681Z","iopub.execute_input":"2022-12-11T07:30:33.088155Z","iopub.status.idle":"2022-12-11T07:30:58.280009Z","shell.execute_reply.started":"2022-12-11T07:30:33.088117Z","shell.execute_reply":"2022-12-11T07:30:58.279104Z"},"trusted":true},"execution_count":null,"outputs":[]},{"cell_type":"markdown","source":"# Co-visitation: click\n- type_x = 0 & type_y = 0\n- delta ts = 1 day","metadata":{}},{"cell_type":"code","source":"%%time\n# USE SMALLEST DISK_PIECES POSSIBLE WITHOUT MEMORY ERROR\nTYPE_COO = \"click\"\nOUTPUT_NAME = f\"top_{TOP_K}_{TYPE_COO}_{DELTA_TS_CLICK}hours\"\nDISK_PIECES = 4 \nSIZE = 1.86e6/DISK_PIECES # total session\n\n# COMPUTE IN PARTS FOR MEMORY MANGEMENT\n\n# 146 files -> 6 CHUNK SIZE 25 -> 1 CHUNK: 5 GROUP SIZE 5\n# 1.86e6 aids -> 4 PART\n\nfor PART in range(DISK_PIECES):\n    print()\n    print('### DISK PART',PART+1)\n    \n    # MERGE IS FASTEST PROCESSING CHUNKS WITHIN CHUNKS\n    # => OUTER CHUNKS\n    for j in range(6):\n        a = j*CHUNK\n        b = min( (j+1)*CHUNK, len(files) )\n        print(f'\\nProcessing files {a} thru {b-1} in groups of {READ_CT}...', end=' ')\n        \n        # => INNER CHUNKS\n        for k in range(a,b,READ_CT):\n            # READ FILE\n            df = [read_file(files[k])]\n            for i in range(1,READ_CT): \n                if k+i<b: df.append( read_file(files[k+i]) )\n            df = cudf.concat(df,ignore_index=True,axis=0)\n            df = df.sort_values(['session','ts'],ascending=[True,False])\n            \n            # USE TAIL OF SESSION\n            df = df.reset_index(drop=True)\n            df['n'] = df.groupby('session').cumcount()\n            df = df.loc[df.n<LENGTH_TAIL_SESSION].drop('n',axis=1)\n            \n            # CREATE PAIRS\n            df = df.merge(df,on='session')\n            df = df.loc[ ((df.ts_x - df.ts_y).abs()< DELTA_TS_CLICK * 60 * 60) & \\\n                                            (df.aid_x != df.aid_y) & \\\n                                            (df.type_y == 0)]\n        \n            # MEMORY MANAGEMENT COMPUTE IN PARTS\n            df = df.loc[(df.aid_x >= PART*SIZE)&(df.aid_x < (PART+1)*SIZE)]\n            \n            # ASSIGN WEIGHTS\n            df = df.drop_duplicates(['session', 'aid_x', 'aid_y'])\n#             df['wgt'] = df.type_y.map(TYPE_WEIGHT_CLICK)\n            df['wgt'] = (df.ts_x - MIN_TS)/(MAX_TS - MIN_TS)\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            # COMBINE INNER CHUNKS\n            if k==a: tmp2 = df\n            else: tmp2 = tmp2.add(df, fill_value=0)\n            print(k,', ',end='')\n        \n        # COMBINE OUTER CHUNKS\n        if a==0: tmp = tmp2\n        else: tmp = tmp.add(tmp2, fill_value=0)\n        del tmp2, df\n        gc.collect()\n        \n        if DEBUG and j >= 0:\n            break\n    # CONVERT MATRIX TO DICTIONARY\n    tmp = tmp.reset_index()\n    tmp = tmp[tmp.wgt >= MIN_CLICK]\n    tmp = tmp.sort_values(['aid_x','wgt'],ascending=[True,False])\n    \n    # SAVE TOP K\n    tmp = tmp.reset_index(drop=True)\n    tmp['n'] = tmp.groupby('aid_x').aid_y.cumcount()\n    tmp = tmp.loc[tmp.n<TOP_K].drop('n',axis=1)\n    tmp = tmp.reset_index(drop=True) # for saving memory\n    \n    # SAVE PART TO DISK (convert to pandas first uses less memory)\n    print(\"Saving sub_parquet...\", end ='')\n    tmp.to_pandas().to_parquet(f'/kaggle/working/{OUTPUT_NAME}_{PART}.pqt')\n    print(\"Saved!\")\n    del tmp\n    gc.collect()\n    if DEBUG and PART >= 0:\n        break\n\nprint(\"Saving full_parquet and full_dict...\", end='')\nfull_parquet = []\nfull_dict = dict()\nfor part in range(DISK_PIECES):\n    sub_parquet = pd.read_parquet(f'/kaggle/working/{OUTPUT_NAME}_{part}.pqt')\n    full_parquet.append(sub_parquet)\n    if SAVE_DICT_FORMAT:\n        full_dict.update(sub_parquet.groupby(\"aid_x\").apply(lambda df: Counter(dict(zip(df.aid_y, df.wgt)))))\n    os.remove(f'/kaggle/working/{OUTPUT_NAME}_{part}.pqt')\n    if DEBUG:\n        break\n    del sub_parquet\n    gc.collect()\nfull_parquet = pd.concat(full_parquet, axis=0, ignore_index=True)\nfull_parquet.to_parquet(f'/kaggle/working/{OUTPUT_NAME}_full.pqt')\nif SAVE_DICT_FORMAT:\n    with open(f'/kaggle/working/{OUTPUT_NAME}_full.pickle', 'wb') as file:\n        pickle.dump(full_dict, file)\nprint(\"Saved!\")\n\ndel full_parquet, full_dict\ngc.collect()","metadata":{"execution":{"iopub.status.busy":"2022-12-05T10:10:25.641085Z","iopub.execute_input":"2022-12-05T10:10:25.64169Z","iopub.status.idle":"2022-12-05T10:10:29.779545Z","shell.execute_reply.started":"2022-12-05T10:10:25.641644Z","shell.execute_reply":"2022-12-05T10:10:29.778521Z"},"trusted":true},"execution_count":null,"outputs":[]},{"cell_type":"markdown","source":"# Co-visitation: cart","metadata":{}},{"cell_type":"code","source":"%%time\nTYPE_COO = \"cart\"\nOUTPUT_NAME = f\"top_{TOP_K}_{TYPE_COO}_{DELTA_TS_BUY}hours\"\n# USE SMALLEST DISK_PIECES POSSIBLE WITHOUT MEMORY ERROR\nDISK_PIECES = 4 \nSIZE = 1.86e6/DISK_PIECES # total session\n\n# COMPUTE IN PARTS FOR MEMORY MANGEMENT\n\n# 146 files -> 6 CHUNK SIZE 25 -> 1 CHUNK: 5 GROUP SIZE 5\n# 1.86e6 aids -> 4 PART\n\nfor PART in range(DISK_PIECES):\n    print()\n    print('### DISK PART',PART+1)\n    \n    # MERGE IS FASTEST PROCESSING CHUNKS WITHIN CHUNKS\n    # => OUTER CHUNKS\n    for j in range(6):\n        a = j*CHUNK\n        b = min( (j+1)*CHUNK, len(files) )\n        print(f'\\nProcessing files {a} thru {b-1} in groups of {READ_CT}...', end=' ')\n        \n        # => INNER CHUNKS\n        for k in range(a,b,READ_CT):\n            # READ FILE\n            df = [read_file(files[k])]\n            for i in range(1,READ_CT): \n                if k+i<b: df.append( read_file(files[k+i]) )\n            df = cudf.concat(df,ignore_index=True,axis=0)\n            df = df.loc[df['type'].isin([0,1])]\n            df = df.sort_values(['session','ts'],ascending=[True,False])\n            \n            # USE TAIL OF SESSION\n            df = df.reset_index(drop=True)\n            df['n'] = df.groupby('session').cumcount()\n            df = df.loc[df.n<LENGTH_TAIL_SESSION].drop('n',axis=1)\n            \n            # CREATE PAIRS\n            df = df.merge(df,on='session')\n            df = df.loc[ ((df.ts_x - df.ts_y).abs()< DELTA_TS_BUY * 60 * 60) & \\\n                                            (df.aid_x != df.aid_y) & \\\n                                            (df.aid_y != 0)]\n        \n            # MEMORY MANAGEMENT COMPUTE IN PARTS\n            df = df.loc[(df.aid_x >= PART*SIZE)&(df.aid_x < (PART+1)*SIZE)]\n            \n            # ASSIGN WEIGHTS\n            df = df.drop_duplicates(['session', 'aid_x', 'aid_y'])\n            df['wgt'] = df.type_y.map(TYPE_WEIGHT_PURCHASE)\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            # COMBINE INNER CHUNKS\n            if k==a: tmp2 = df\n            else: tmp2 = tmp2.add(df, fill_value=0)\n            print(k,', ',end='')\n        \n        # COMBINE OUTER CHUNKS\n        if a==0: tmp = tmp2\n        else: tmp = tmp.add(tmp2, fill_value=0)\n        del tmp2, df\n        gc.collect()\n        \n        if DEBUG and j >= 0:\n            break\n            \n    # CONVERT MATRIX TO DICTIONARY\n    tmp = tmp.reset_index()\n    tmp = tmp[tmp.wgt >= MIN_CO]\n    tmp = tmp.sort_values(['aid_x','wgt'],ascending=[True,False])\n\n    # SAVE TOP K\n    tmp = tmp.reset_index(drop=True)\n    tmp['n'] = tmp.groupby('aid_x').aid_y.cumcount()\n    tmp = tmp.loc[tmp.n<TOP_K].drop('n',axis=1)\n    tmp = tmp.reset_index(drop=True) # for saving memory\n\n    # SAVE PART TO DISK (convert to pandas first uses less memory)\n    print(\"Saving sub_parquet...\", end ='')\n    tmp.to_pandas().to_parquet(f'/kaggle/working/{OUTPUT_NAME}_{PART}.pqt')\n    print(\"Saved!\")\n    del tmp\n    gc.collect()\n    if DEBUG and PART >= 0:\n        break\n\nprint(\"Saving full_parquet and full_dict...\", end='')\nfull_parquet = []\nfull_dict = dict()\nfor part in range(DISK_PIECES):\n    sub_parquet = pd.read_parquet(f'/kaggle/working/{OUTPUT_NAME}_{part}.pqt')\n    full_parquet.append(sub_parquet)\n    if SAVE_DICT_FORMAT:\n        full_dict.update(sub_parquet.groupby(\"aid_x\").apply(lambda df: Counter(dict(zip(df.aid_y, df.wgt)))))\n    os.remove(f'/kaggle/working/{OUTPUT_NAME}_{part}.pqt')\n    if DEBUG:\n        break\n    del sub_parquet\n    gc.collect()\nfull_parquet = pd.concat(full_parquet, axis=0, ignore_index=True)\nfull_parquet.to_parquet(f'/kaggle/working/{OUTPUT_NAME}_full.pqt')\nif SAVE_DICT_FORMAT:\n    with open(f'/kaggle/working/{OUTPUT_NAME}_full.pickle', 'wb') as file:\n        pickle.dump(full_dict, file)\nprint(\"Saved!\")\n\ndel full_parquet, full_dict\ngc.collect()","metadata":{},"execution_count":null,"outputs":[]},{"cell_type":"markdown","source":"# Co-visitation: purchase\n> - delta ts = 1 day","metadata":{}},{"cell_type":"code","source":"%%time\nTYPE_COO = \"purchase\"\nOUTPUT_NAME = f\"top_{TOP_K}_{TYPE_COO}_{DELTA_TS_BUY}hours\"\n# USE SMALLEST DISK_PIECES POSSIBLE WITHOUT MEMORY ERROR\nDISK_PIECES = 4 \nSIZE = 1.86e6/DISK_PIECES # total session\n\n# COMPUTE IN PARTS FOR MEMORY MANGEMENT\n\n# 146 files -> 6 CHUNK SIZE 25 -> 1 CHUNK: 5 GROUP SIZE 5\n# 1.86e6 aids -> 4 PART\n\nfor PART in range(DISK_PIECES):\n    print()\n    print('### DISK PART',PART+1)\n    \n    # MERGE IS FASTEST PROCESSING CHUNKS WITHIN CHUNKS\n    # => OUTER CHUNKS\n    for j in range(6):\n        a = j*CHUNK\n        b = min( (j+1)*CHUNK, len(files) )\n        print(f'\\nProcessing files {a} thru {b-1} in groups of {READ_CT}...', end=' ')\n        \n        # => INNER CHUNKS\n        for k in range(a,b,READ_CT):\n            # READ FILE\n            df = [read_file(files[k])]\n            for i in range(1,READ_CT): \n                if k+i<b: df.append( read_file(files[k+i]) )\n            df = cudf.concat(df,ignore_index=True,axis=0)\n            df = df.sort_values(['session','ts'],ascending=[True,False])\n            \n            # USE TAIL OF SESSION\n            df = df.reset_index(drop=True)\n            df['n'] = df.groupby('session').cumcount()\n            df = df.loc[df.n<LENGTH_TAIL_SESSION].drop('n',axis=1)\n            \n            # CREATE PAIRS\n            df = df.merge(df,on='session')\n            df = df.loc[ ((df.ts_x - df.ts_y).abs()< DELTA_TS_BUY * 60 * 60) & \\\n                                            (df.aid_x != df.aid_y) & \\\n                                            (df.type_y != 0) & (df.type_y != 1)]\n        \n            # MEMORY MANAGEMENT COMPUTE IN PARTS\n            df = df.loc[(df.aid_x >= PART*SIZE)&(df.aid_x < (PART+1)*SIZE)]\n            \n            # ASSIGN WEIGHTS\n            df = df.drop_duplicates(['session', 'aid_x', 'aid_y'])\n            df['wgt'] = df.type_y.map(TYPE_WEIGHT_PURCHASE)\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            # COMBINE INNER CHUNKS\n            if k==a: tmp2 = df\n            else: tmp2 = tmp2.add(df, fill_value=0)\n            print(k,', ',end='')\n        \n        # COMBINE OUTER CHUNKS\n        if a==0: tmp = tmp2\n        else: tmp = tmp.add(tmp2, fill_value=0)\n        del tmp2, df\n        gc.collect()\n        \n        if DEBUG and j >= 0:\n            break\n            \n    # CONVERT MATRIX TO DICTIONARY\n    tmp = tmp.reset_index()\n    tmp = tmp[tmp.wgt >= MIN_CO]\n    tmp = tmp.sort_values(['aid_x','wgt'],ascending=[True,False])\n\n    # SAVE TOP K\n    tmp = tmp.reset_index(drop=True)\n    tmp['n'] = tmp.groupby('aid_x').aid_y.cumcount()\n    tmp = tmp.loc[tmp.n<TOP_K].drop('n',axis=1)\n    tmp = tmp.reset_index(drop=True) # for saving memory\n\n    # SAVE PART TO DISK (convert to pandas first uses less memory)\n    print(\"Saving sub_parquet...\", end ='')\n    tmp.to_pandas().to_parquet(f'/kaggle/working/{OUTPUT_NAME}_{PART}.pqt')\n    print(\"Saved!\")\n    del tmp\n    gc.collect()\n    if DEBUG and PART >= 0:\n        break\n\nprint(\"Saving full_parquet and full_dict...\", end='')\nfull_parquet = []\nfull_dict = dict()\nfor part in range(DISK_PIECES):\n    sub_parquet = pd.read_parquet(f'/kaggle/working/{OUTPUT_NAME}_{part}.pqt')\n    full_parquet.append(sub_parquet)\n    if SAVE_DICT_FORMAT:\n        full_dict.update(sub_parquet.groupby(\"aid_x\").apply(lambda df: Counter(dict(zip(df.aid_y, df.wgt)))))\n    os.remove(f'/kaggle/working/{OUTPUT_NAME}_{part}.pqt')\n    if DEBUG:\n        break\n    del sub_parquet\n    gc.collect()\nfull_parquet = pd.concat(full_parquet, axis=0, ignore_index=True)\nfull_parquet.to_parquet(f'/kaggle/working/{OUTPUT_NAME}_full.pqt')\nif SAVE_DICT_FORMAT:\n    with open(f'/kaggle/working/{OUTPUT_NAME}_full.pickle', 'wb') as file:\n        pickle.dump(full_dict, file)\nprint(\"Saved!\")\n\ndel full_parquet, full_dict\ngc.collect()","metadata":{"execution":{"iopub.status.busy":"2022-12-04T02:16:33.806172Z","iopub.execute_input":"2022-12-04T02:16:33.806569Z","iopub.status.idle":"2022-12-04T02:16:36.132038Z","shell.execute_reply.started":"2022-12-04T02:16:33.80653Z","shell.execute_reply":"2022-12-04T02:16:36.130186Z"},"trusted":true},"execution_count":null,"outputs":[]}]}