{"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 is a fork of carnozhao's [notebook](https://www.kaggle.com/code/carnozhao/otto-fast-cpu-end-to-end-pipeline).❗❗❗** \n\nOverall the code used is not too different from the original but I tuned the various parameters inside it.\n\nAn important key to improving the quality of the solution was reducing a lot the time frame for co-visitation and using a different time frame for both clicks and carts/orders. \n\nI reduced from 24 hours to 0.06 and 0.07 hours, so the best time frame was just some minutes.\n\n\n# **❗❗❗This is a Work in Progress, I'll try to make it more readable in the weekend. If you have doubts, or you want more information about something, feel free to comment and I'll try to answer as soon as I can❗❗❗**","metadata":{}},{"cell_type":"code","source":"import os\nimport gc\nimport heapq\nimport pickle\nimport numba as nb\nimport numpy as np\nimport pandas as pd\nfrom tqdm.auto import tqdm","metadata":{"_uuid":"8f2839f25d086af736a60e9eeb907d3b93b6e0e5","_cell_guid":"b1076dfc-b9ad-4769-8c92-a6c4dae69d19","execution":{"iopub.status.busy":"2022-11-29T22:39:21.138292Z","iopub.execute_input":"2022-11-29T22:39:21.138731Z","iopub.status.idle":"2022-11-29T22:39:22.139968Z","shell.execute_reply.started":"2022-11-29T22:39:21.138655Z","shell.execute_reply":"2022-11-29T22:39:22.138734Z"},"trusted":true},"execution_count":null,"outputs":[]},{"cell_type":"markdown","source":"### Data\n\n#### Meta DataFrame\n\n`train.csv/test.csv` columns:\n\n- `session`: ordered session ids\n- `length`: length of each session\n- `start_time`: start time stamp (unit=1sec) of this session\n\n#### Data Array\n\n`train.npz/test.npz`:\n\nAn npz contains three 1-dim arrays: \n\n- `aids`: aid array\n- `ts`: time stamp array (minused by start_time of corresponding session)\n- `ops`: type array, `clicks=0, carts=1, orders=2`\n\nAn array is the concatenation of each session, ordered by session id\n\ne.g.\n\n```python\naids = np.concatenate([\n        [1, 2, 3], # session 0 time-ordered aids\n        [4, 5, 6], # session 1 time-ordered aids\n        ...])\n```\n","metadata":{}},{"cell_type":"code","source":"params={\n    \"ops_weights\" : np.array([1.0, 6.0, 3.0]),\n    \"test_ops_weights\" : np.array([1.0, 6.0, 3.0]),\n    \"time_weight\":[3,3,3],\n    \"hours_co_visitation_click\":0.07,\n    \"hours_co_visitation_buy\":0.06,\n    \"topn\":20,\n    \"seq_weight\":[-3,-4,-2], \n    \"rank_w\":[-4,-4,-5],\n    \"base_time_weight\":[1.4,1.4,1.4],\n    \"tail\":100\n}\n\ntail = params[\"tail\"]\nparallel = 1024\n\n\nOP_WEIGHT = 0; TIME_WEIGHT = 1\n\nops_weights = params[\"ops_weights\"]\ntest_ops_weights = params[\"test_ops_weights\"]\ntopn = params[\"topn\"]\n\n\n\ntime_covisitation=[0,0]\ntime_covisitation[TIME_WEIGHT]=params[\"hours_co_visitation_click\"]\ntime_covisitation[OP_WEIGHT]=params[\"hours_co_visitation_buy\"]","metadata":{},"execution_count":null,"outputs":[]},{"cell_type":"code","source":"df = pd.read_csv(\"../input/otto-data/train.csv\")\ndf_test = pd.read_csv(\"../input/otto-data/test.csv\")\ndf = pd.concat([df, df_test]).reset_index(drop = True)\nnpz = np.load(\"../input/otto-data/train.npz\")\nnpz_test = np.load(\"../input/otto-data/test.npz\")\naids = np.concatenate([npz['aids'], npz_test['aids']])\nts = np.concatenate([npz['ts'], npz_test['ts']])\nops = np.concatenate([npz['ops'], npz_test['ops']])\n\ndf[\"idx\"] = np.cumsum(df.length) - df.length\ndf[\"end_time\"] = df.start_time + ts[df.idx + df.length - 1]","metadata":{"execution":{"iopub.status.busy":"2022-11-29T22:39:22.141721Z","iopub.execute_input":"2022-11-29T22:39:22.142024Z","iopub.status.idle":"2022-11-29T22:39:56.324244Z","shell.execute_reply.started":"2022-11-29T22:39:22.141993Z","shell.execute_reply":"2022-11-29T22:39:56.322425Z"},"trusted":true},"execution_count":null,"outputs":[]},{"cell_type":"code","source":"len_test=len(df_test)","metadata":{},"execution_count":null,"outputs":[]},{"cell_type":"code","source":"aids.dtype,ts.dtype,ops.dtype","metadata":{"execution":{"iopub.status.busy":"2022-11-29T22:39:56.326234Z","iopub.execute_input":"2022-11-29T22:39:56.326749Z","iopub.status.idle":"2022-11-29T22:39:56.337452Z","shell.execute_reply.started":"2022-11-29T22:39:56.326716Z","shell.execute_reply":"2022-11-29T22:39:56.336489Z"},"trusted":true},"execution_count":null,"outputs":[]},{"cell_type":"code","source":"aids[:3],ts[:3],ops[:3]","metadata":{"execution":{"iopub.status.busy":"2022-11-29T22:42:33.654848Z","iopub.execute_input":"2022-11-29T22:42:33.655175Z","iopub.status.idle":"2022-11-29T22:42:33.666398Z","shell.execute_reply.started":"2022-11-29T22:42:33.655151Z","shell.execute_reply":"2022-11-29T22:42:33.665097Z"},"trusted":true},"execution_count":null,"outputs":[]},{"cell_type":"markdown","source":"### Numba functions\n\nTo fully utilize the ability of numba's jit, and to avoid memory errors, we process `parallel=1024` sessions in one jit function call.\n\n- Firstly, we count the `(aid1, aid2)` pairs in each session in parallel. \n\n- Then, we merge these pairs into a nested counter `{aid1: {aid2: n}, ...}`\n\n- Finally, we find the top-k `aid2` from `aid1`'s counter","metadata":{}},{"cell_type":"code","source":"# get pair dict {(aid1, aid2): weight} for each session\n# The maximum time span between two points is 1 day = 24 * 60 * 60 sec\n@nb.jit(nopython = True, cache = True)\ndef get_single_pairs(pairs, aids, ts, ops, idx, length, start_time, ops_weights, mode,time_weight,hours_co_visitation,base_time_weight):\n    max_idx = idx + length\n    min_idx = max(max_idx - tail, idx)\n    for i in range(min_idx, max_idx):\n        for j in range(i + 1, max_idx):\n            if ts[j] - ts[i] >= hours_co_visitation * 60 * 60: break\n            if aids[i] == aids[j]: continue\n            if mode == OP_WEIGHT:\n                w1 = ops_weights[ops[j]]\n                w2 = ops_weights[ops[i]]\n            elif mode == TIME_WEIGHT:\n                w1 = base_time_weight[ops[j]] + time_weight[ops[j]] * (ts[i] + start_time - 1659304800) / (1662328791 - 1659304800)\n                w2 = base_time_weight[ops[i]] + time_weight[ops[j]] * (ts[j] + start_time - 1659304800) / (1662328791 - 1659304800)\n            pairs[(aids[i], aids[j])] = w1\n            pairs[(aids[j], aids[i])] = w2\n\n# get pair dict of each session in parallel\n# merge pairs into a nested dict format (cnt)\n@nb.jit(nopython = True, parallel = True, cache = True)\ndef get_pairs(aids, ts, ops, row, cnts, ops_weights, mode,time_weight,hours_co_visitation,base_time_weight):\n    par_n = len(row)\n    pairs = [{(0, 0): 0.0 for _ in range(0)} for _ in range(par_n)]\n    for par_i in nb.prange(par_n):\n        _, idx, length, start_time = row[par_i]\n        get_single_pairs(pairs[par_i], aids, ts, ops, idx, length, start_time, ops_weights, mode,time_weight,hours_co_visitation,base_time_weight)\n    for par_i in range(par_n):\n        for (aid1, aid2), w in pairs[par_i].items():\n            if aid1 not in cnts: cnts[aid1] = {0: 0.0 for _ in range(0)}\n            cnt = cnts[aid1]\n            if aid2 not in cnt: cnt[aid2] = 0.0\n            cnt[aid2] += w\n    \n# util function to get most common keys from a counter dict using min-heap\n# overwrite == 1 means the later item with equal weight is more important\n# otherwise, means the former item with equal weight is more important\n# the result is ordered from higher weight to lower weight\n@nb.jit(nopython = True, cache = True)\ndef heap_topk(cnt, overwrite, cap):\n    q = [(0.0, 0, 0) for _ in range(0)]\n    for i, (k, n) in enumerate(cnt.items()):\n        if overwrite == 1:\n            heapq.heappush(q, (n, i, k))\n        else:\n            heapq.heappush(q, (n, -i, k))\n        if len(q) > cap:\n            heapq.heappop(q)\n    return [heapq.heappop(q)[2] for _ in range(len(q))][::-1]\n   \n# save top-k aid2 for each aid1's cnt\n@nb.jit(nopython = True, cache = True)\ndef get_topk(cnts, topk, k):\n    for aid1, cnt in cnts.items():\n        topk[aid1] = np.array(heap_topk(cnt, 1, k))","metadata":{"execution":{"iopub.status.busy":"2022-11-13T14:40:54.209284Z","iopub.execute_input":"2022-11-13T14:40:54.209692Z","iopub.status.idle":"2022-11-13T14:40:54.471035Z","shell.execute_reply.started":"2022-11-13T14:40:54.209657Z","shell.execute_reply":"2022-11-13T14:40:54.469904Z"},"trusted":true},"execution_count":null,"outputs":[]},{"cell_type":"markdown","source":"### Train\n\nHere we use topk counting from [Deotte's notebooks](https://www.kaggle.com/code/cdeotte/candidate-rerank-model-lb-0-573). I dropped some of his ideas to make the whole pipeline simple but still strong.\n\n- **mode 0**: Counters are weighted by operation types, with `op_weights = [1.0, 6.0, 3.0]`. This will be used for carts and orders prediction.\n\n- **mode 1**: Counters are weighted by operation time.  This will be used for clicks prediction.\n\nFor each mode, the eastimated running time is **7min~8min** in kaggle notebook","metadata":{}},{"cell_type":"code","source":"topks = {}\n\n# for two modes\nfor mode in [OP_WEIGHT, TIME_WEIGHT]:\n    # get nested counter\n    cnts = nb.typed.Dict.empty(\n        key_type = nb.types.int64,\n        value_type = nb.typeof(nb.typed.Dict.empty(key_type = nb.types.int64, value_type = nb.types.float64)))\n    max_idx = len(df)\n    for idx in tqdm(range(0, max_idx, parallel)):\n        row = df.iloc[idx:min(idx + parallel, max_idx)][['session', 'idx', 'length', 'start_time']].values\n        get_pairs(aids, ts, ops, row, cnts, ops_weights, mode,params[\"time_weight\"],time_covisitation[mode],params[\"base_time_weight\"])\n\n    # get topk from counter\n    topk = nb.typed.Dict.empty(\n            key_type = nb.types.int64,\n            value_type = nb.types.int64[:])\n    get_topk(cnts, topk, topn)\n\n    del cnts; gc.collect()\n    topks[mode] = topk\ngc.collect()","metadata":{"execution":{"iopub.status.busy":"2022-11-13T14:40:54.474038Z","iopub.execute_input":"2022-11-13T14:40:54.474491Z","iopub.status.idle":"2022-11-13T14:57:30.234723Z","shell.execute_reply.started":"2022-11-13T14:40:54.474443Z","shell.execute_reply":"2022-11-13T14:57:30.233456Z"},"trusted":true},"execution_count":null,"outputs":[]},{"cell_type":"markdown","source":"### Inference Numba Functions\n\nThe logic is a simpler version of [Deotte's notebook](https://www.kaggle.com/code/cdeotte/candidate-rerank-model-lb-0-573).\n\nSame as training functions, we process `parallel=1024` sessions in one jit function call. For each session, if the unique aids is more than 20, we rerank them by a time-decay weight. Otherwise, we use current aids to recall new aids from their topk co-visitation candidates, weighted by the query's operation type (`test_ops_weights=[1.0, 6.0, 3.0]`). \n\nWe used **mode0** topk to generate predictions for carts and orders, and **mode1** topk to generate predictions for clicks. ","metadata":{}},{"cell_type":"code","source":"@nb.jit(nopython = True, cache = True)\ndef inference_(aids, ops, row, result, topk, test_ops_weights, seq_weight, rank_w):\n    for session, idx, length in row:\n        rank_weight = np.power(2, np.linspace(rank_w, 0, topn))[::-1]\n        unique_aids = nb.typed.Dict.empty(key_type = nb.types.int64, value_type = nb.types.float64)\n        cnt = nb.typed.Dict.empty(key_type = nb.types.int64, value_type = nb.types.float64)\n        cnt_additional_items = nb.typed.Dict.empty(key_type = nb.types.int64, value_type = nb.types.float64)\n        \n        candidates = aids[idx:idx + length][::-1]\n        candidates_ops = ops[idx:idx + length][::-1]\n        for a in candidates:\n            unique_aids[a] = 0\n        sequence_weight = np.power(2, np.linspace(seq_weight, 0, len(candidates)))[::-1]\n        for a, op, w in zip(candidates, candidates_ops, sequence_weight):\n            if a not in cnt: \n                cnt[a] = 0\n            cnt[a] += w * test_ops_weights[op]\n            \n        if len(unique_aids) >= 20:\n            result_candidates = heap_topk(cnt, 0, 20)\n        else:\n            result_candidates = list(unique_aids)\n            for a,w in zip(cnt,sequence_weight):\n                if a not in topk: \n                    continue\n                for rank, b in enumerate(topk[a]):\n                    if b in unique_aids: \n                        continue\n                    if b not in cnt_additional_items: \n                        cnt_additional_items[b] = 0\n                    cnt_additional_items[b] += w*rank_weight[rank] # add weight equal to the frequency of the weight of the corresponding item\n            result_candidates.extend(heap_topk(cnt_additional_items, 0, 20 - len(result_candidates)))\n        result[session] = np.array(result_candidates)\n        \n\n@nb.jit(nopython = True)\ndef inference(aids, ops, row, \n              result_clicks, result_buy,result_order,\n              topk_clicks, topk_buy,\n              test_ops_weights,seq_weight,rank_w):\n    inference_(aids, ops, row, result_clicks, topk_clicks, test_ops_weights, seq_weight[0],rank_w[0])\n    inference_(aids, ops, row, result_buy, topk_buy, test_ops_weights, seq_weight[1],rank_w[1])\n    inference_(aids, ops, row, result_order, topk_buy, test_ops_weights, seq_weight[2],rank_w[2])","metadata":{"execution":{"iopub.status.busy":"2022-11-13T15:01:48.75753Z","iopub.execute_input":"2022-11-13T15:01:48.757954Z","iopub.status.idle":"2022-11-13T15:01:48.78381Z","shell.execute_reply.started":"2022-11-13T15:01:48.757921Z","shell.execute_reply":"2022-11-13T15:01:48.782726Z"},"trusted":true},"execution_count":null,"outputs":[]},{"cell_type":"code","source":"# result place holder\nresult_clicks = nb.typed.Dict.empty(\n    key_type = nb.types.int64,\n    value_type = nb.types.int64[:])\nresult_buy = nb.typed.Dict.empty(\n    key_type = nb.types.int64,\n    value_type = nb.types.int64[:])\nresult_order = nb.typed.Dict.empty(\n    key_type = nb.types.int64,\n    value_type = nb.types.int64[:])\nfor idx in tqdm(range(len(df) - len_test, len(df), parallel)):\n    row = df.iloc[idx:min(idx + parallel, len(df))][['session', 'idx', 'length']].values\n    inference(aids, ops, row, result_clicks, result_buy,result_order, topks[TIME_WEIGHT], topks[OP_WEIGHT], test_ops_weights,params[\"seq_weight\"],params[\"rank_w\"])","metadata":{"execution":{"iopub.status.busy":"2022-11-13T15:01:49.54435Z","iopub.execute_input":"2022-11-13T15:01:49.544781Z","iopub.status.idle":"2022-11-13T15:01:51.991856Z","shell.execute_reply.started":"2022-11-13T15:01:49.544743Z","shell.execute_reply":"2022-11-13T15:01:51.990734Z"},"trusted":true},"execution_count":null,"outputs":[]},{"cell_type":"code","source":"subs = []\nop_names = [\"clicks\", \"carts\", \"orders\"]\nfor result, op in zip([result_clicks, result_buy, result_order], op_names):\n\n    sub = pd.DataFrame({\"session_type\": result.keys(), \"labels\": result.values()})\n    sub.session_type = sub.session_type.astype(str) + f\"_{op}\"\n    sub.labels = sub.labels.apply(lambda x: \" \".join(x.astype(str)))\n    subs.append(sub)\n    \nsub = pd.concat(subs).reset_index(drop = True)\nsub.to_csv('submission.csv', index = False)\nsub.head()","metadata":{"execution":{"iopub.status.busy":"2022-11-13T15:01:51.993623Z","iopub.execute_input":"2022-11-13T15:01:51.99398Z","iopub.status.idle":"2022-11-13T15:01:52.756408Z","shell.execute_reply.started":"2022-11-13T15:01:51.993947Z","shell.execute_reply":"2022-11-13T15:01:52.755261Z"},"trusted":true},"execution_count":null,"outputs":[]}]}