{"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"},"kaggle":{"accelerator":"gpu","dataSources":[{"sourceId":38760,"databundleVersionId":4493939,"sourceType":"competition"},{"sourceId":4436180,"sourceType":"datasetVersion","datasetId":2597726}],"dockerImageVersionId":30302,"isInternetEnabled":true,"language":"python","sourceType":"notebook","isGpuEnabled":true}},"nbformat_minor":4,"nbformat":4,"cells":[{"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","metadata":{"papermill":{"duration":3.036143,"end_time":"2022-11-10T16:03:24.014816","exception":false,"start_time":"2022-11-10T16:03:20.978673","status":"completed"},"tags":[],"execution":{"iopub.status.busy":"2022-11-16T17:59:30.852609Z","iopub.execute_input":"2022-11-16T17:59:30.853062Z","iopub.status.idle":"2022-11-16T17:59:33.538931Z","shell.execute_reply.started":"2022-11-16T17:59:30.852976Z","shell.execute_reply":"2022-11-16T17:59:33.537552Z"},"trusted":true},"outputs":[],"execution_count":null},{"cell_type":"code","source":"# Cache functions\ndef read_file(f):\n    return cudf.DataFrame( data_cache[f] )\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 dữ liệu lên CPU\ndata_cache = {}\ntype_labels = {'clicks':0, 'carts':1, 'orders':2}\nfiles = glob.glob('../input/otto-chunk-data-inparquet-format/*_parquet/*')\nfor f in files: data_cache[f] = read_file_to_cache(f)\n\n# Tham số chunk\nREAD_CT = 5\nCHUNK = int( np.ceil( len(files)/6 ))","metadata":{"papermill":{"duration":0.063943,"end_time":"2022-11-10T16:03:24.091816","exception":false,"start_time":"2022-11-10T16:03:24.027873","status":"completed"},"tags":[],"_kg_hide-input":true,"execution":{"iopub.status.busy":"2022-11-16T17:59:33.542845Z","iopub.execute_input":"2022-11-16T17:59:33.543142Z","iopub.status.idle":"2022-11-16T18:00:35.829684Z","shell.execute_reply.started":"2022-11-16T17:59:33.543116Z","shell.execute_reply":"2022-11-16T18:00:35.828024Z"},"trusted":true},"outputs":[],"execution_count":null},{"cell_type":"code","source":"type_weight = {0:1, 1:6, 2:3}\n\nDISK_PIECES = 4\nSIZE = 1.86e6/DISK_PIECES\n\n# Tính toán theo từng phần\nfor PART in range(DISK_PIECES):\n    print()\n    print('### DISK PART',PART+1)\n    \n    # Outer chunks\n    for j in range(6):\n        a = j*CHUNK\n        b = min( (j+1)*CHUNK, len(files) )\n        print(f'Processing files {a} thru {b-1} in groups of {READ_CT}...')\n        \n        # Inner chunks\n        for k in range(a,b,READ_CT):\n            # Đọc 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            # Chỉ dùng chuỗi 30 sự kiện cuối cùng\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            # Tạo cặp\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            df = df.loc[(df.aid_x >= PART*SIZE)&(df.aid_x < (PART+1)*SIZE)]\n            # Gán trọng số\n            df = df[['session', 'aid_x', 'aid_y','type_y']].drop_duplicates(['session', 'aid_x', 'aid_y'])\n            df['wgt'] = df.type_y.map(type_weight)\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            # Gộp inner chunks\n            if k==a: tmp2 = df\n            else: tmp2 = tmp2.add(df, fill_value=0)\n            print(k,', ',end='')\n        print()\n        # Gộp 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    # Chuyển ma trận về dictionary\n    tmp = tmp.reset_index()\n    tmp = tmp.sort_values(['aid_x','wgt'],ascending=[True,False])\n    # Lưu top15\n    tmp = tmp.reset_index(drop=True)\n    tmp['n'] = tmp.groupby('aid_x').aid_y.cumcount()\n    tmp = tmp.loc[tmp.n<15].drop('n',axis=1)\n    # Lưu lại\n    tmp.to_pandas().to_parquet(f'top_15_carts_orders_{PART}.pqt')","metadata":{"papermill":{"duration":566.561189,"end_time":"2022-11-10T16:12:50.666123","exception":false,"start_time":"2022-11-10T16:03:24.104934","status":"completed"},"tags":[],"execution":{"iopub.status.busy":"2022-11-16T18:00:35.83222Z","iopub.execute_input":"2022-11-16T18:00:35.832992Z","iopub.status.idle":"2022-11-16T18:03:55.366191Z","shell.execute_reply.started":"2022-11-16T18:00:35.832947Z","shell.execute_reply":"2022-11-16T18:03:55.365039Z"},"trusted":true},"outputs":[],"execution_count":null},{"cell_type":"code","source":"DISK_PIECES = 1\nSIZE = 1.86e6/DISK_PIECES\n\n# Tính toán theo từng phần\nfor PART in range(DISK_PIECES):\n    print()\n    print('### DISK PART',PART+1)\n    \n    # Outer chunks\n    for j in range(6):\n        a = j*CHUNK\n        b = min( (j+1)*CHUNK, len(files) )\n        print(f'Processing files {a} thru {b-1} in groups of {READ_CT}...')\n        \n        # Inner chunk\n        for k in range(a,b,READ_CT):\n            # Đọc 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([1,2])] # ONLY WANT CARTS AND ORDERS\n            df = df.sort_values(['session','ts'],ascending=[True,False])\n            # Dùng chuỗi 30 sự kiện cuối cùng\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            # Tạo cặp\n            df = df.merge(df,on='session')\n            df = df.loc[ ((df.ts_x - df.ts_y).abs()< 14 * 24 * 60 * 60) & (df.aid_x != df.aid_y) ] # 14 DAYS\n            df = df.loc[(df.aid_x >= PART*SIZE)&(df.aid_x < (PART+1)*SIZE)]\n            # Gán trọng số\n            df = df[['session', 'aid_x', 'aid_y','type_y']].drop_duplicates(['session', 'aid_x', 'aid_y'])\n            df['wgt'] = 1\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            # Gộp inner chunks\n            if k==a: tmp2 = df\n            else: tmp2 = tmp2.add(df, fill_value=0)\n            print(k,', ',end='')\n        print()\n        # Gộp 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    # Chuyển ma trận thành dictionary\n    tmp = tmp.reset_index()\n    tmp = tmp.sort_values(['aid_x','wgt'],ascending=[True,False])\n    # Lưu top15\n    tmp = tmp.reset_index(drop=True)\n    tmp['n'] = tmp.groupby('aid_x').aid_y.cumcount()\n    tmp = tmp.loc[tmp.n<15].drop('n',axis=1)\n    # Lưu lại\n    tmp.to_pandas().to_parquet(f'top_15_buy2buy_{PART}.pqt')","metadata":{"_kg_hide-input":true,"_kg_hide-output":true,"papermill":{"duration":113.735315,"end_time":"2022-11-10T16:14:44.498182","exception":false,"start_time":"2022-11-10T16:12:50.762867","status":"completed"},"tags":[],"execution":{"iopub.status.busy":"2022-11-16T18:03:55.369013Z","iopub.execute_input":"2022-11-16T18:03:55.369394Z","iopub.status.idle":"2022-11-16T18:04:26.289606Z","shell.execute_reply.started":"2022-11-16T18:03:55.369354Z","shell.execute_reply":"2022-11-16T18:04:26.288386Z"},"trusted":true},"outputs":[],"execution_count":null},{"cell_type":"code","source":"DISK_PIECES = 4\nSIZE = 1.86e6/DISK_PIECES\n\n# Tính toán theo từng phần\nfor PART in range(DISK_PIECES):\n    print()\n    print('### DISK PART',PART+1)\n    \n    # Outer chunks\n    for j in range(6):\n        a = j*CHUNK\n        b = min( (j+1)*CHUNK, len(files) )\n        print(f'Processing files {a} thru {b-1} in groups of {READ_CT}...')\n        \n        # Inner chunks\n        for k in range(a,b,READ_CT):\n            # Đọc 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            # Dùng chuỗi 30 sự kiện cuối cùng\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            # Tạo cặp\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            df = df.loc[(df.aid_x >= PART*SIZE)&(df.aid_x < (PART+1)*SIZE)]\n            # Gán trọng số\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            # Gộp inner chunks\n            if k==a: tmp2 = df\n            else: tmp2 = tmp2.add(df, fill_value=0)\n            print(k,', ',end='')\n        print()\n        # Gộp 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    # Chuyển ma trận về dictionary\n    tmp = tmp.reset_index()\n    tmp = tmp.sort_values(['aid_x','wgt'],ascending=[True,False])\n    # Lưu top20\n    tmp = tmp.reset_index(drop=True)\n    tmp['n'] = tmp.groupby('aid_x').aid_y.cumcount()\n    tmp = tmp.loc[tmp.n<20].drop('n',axis=1)\n    # Lưu lại\n    tmp.to_pandas().to_parquet(f'top_20_clicks_{PART}.pqt')","metadata":{"_kg_hide-input":true,"_kg_hide-output":true,"papermill":{"duration":null,"end_time":null,"exception":false,"start_time":"2022-11-10T16:14:44.629032","status":"running"},"tags":[],"execution":{"iopub.status.busy":"2022-11-16T18:04:26.291394Z","iopub.execute_input":"2022-11-16T18:04:26.291844Z","iopub.status.idle":"2022-11-16T18:07:42.542742Z","shell.execute_reply.started":"2022-11-16T18:04:26.291801Z","shell.execute_reply":"2022-11-16T18:07:42.54165Z"},"trusted":true},"outputs":[],"execution_count":null},{"cell_type":"code","source":"del data_cache, tmp\n_ = gc.collect()","metadata":{"execution":{"iopub.status.busy":"2022-11-16T18:07:42.544185Z","iopub.execute_input":"2022-11-16T18:07:42.54497Z","iopub.status.idle":"2022-11-16T18:07:42.700407Z","shell.execute_reply.started":"2022-11-16T18:07:42.54493Z","shell.execute_reply":"2022-11-16T18:07:42.699305Z"},"trusted":true},"outputs":[],"execution_count":null}]}