{"cells":[{"cell_type":"markdown","metadata":{},"source":"# OTTO – Multi-Objective Recommender System\n*Late-submission solution by **Mohamed Elshetiwi** (no1mohamed)*\n\nThis notebook implements a covisitation-matrix candidate generator + per-session\nranker. It produces a `submission.csv` directly from this notebook environment\nusing the official competition test set (in parquet form, via Radek's\n`otto-full-optimized-memory-footprint` dataset).\n\n**Architecture**\n1. Build three time-decayed covisitation matrices: click→click, click→buy, buy→buy\n2. Generate per-session candidate lists (top-100) as union of:\n   own-session items + covis neighbours + popularity fallback\n3. Take top-20 per (session, target_type), format as Kaggle submission\n\nTargets: clicks, carts, orders. Metric: 0.10·clk_R@20 + 0.30·cart_R@20 + 0.60·ord_R@20\n\nLocal validation uses Radek's `otto-train-and-test-data-for-local-validation` split\n(the last 7 days of train held out exactly as the public LB does).\n"},{"cell_type":"code","metadata":{},"execution_count":null,"outputs":[],"source":"# --- imports & paths ---\nimport os, gc, time, json\nfrom pathlib import Path\nfrom collections import defaultdict\nimport numpy as np\nimport polars as pl\nprint('polars', pl.__version__)\n\n# Datasets attached on the right:\nDATA_FULL  = Path('/kaggle/input/otto-full-optimized-memory-footprint')   # full train + official test, parquet\nDATA_RADEK = Path('/kaggle/input/otto-train-and-test-data-for-local-validation')  # local-validation split\n\nOUT = Path('/kaggle/working'); OUT.mkdir(exist_ok=True)\nART = OUT / 'artifacts'; ART.mkdir(exist_ok=True)\n\nprint('FULL exists?', DATA_FULL.exists(), '   files:', sorted(p.name for p in DATA_FULL.glob('*')) if DATA_FULL.exists() else 'NO DIR')\nprint('RADEK exists?', DATA_RADEK.exists(), '  files:', sorted(p.name for p in DATA_RADEK.glob('*')) if DATA_RADEK.exists() else 'NO DIR')\n"},{"cell_type":"markdown","metadata":{},"source":"## 1. Build covisitation matrices\n\nThree matrices, each (aid_x → list of top-K aid_y, weight). We use Radek's\npre-converted train.parquet (the local-validation training pool — 11M sessions,\n~164M events, ts in seconds).\n\nFor each session: keep last `N_TAIL=30` events. Generate ordered pairs whose\ntimestamps are within 24 hours. Pairs are weighted: `click2buy` weights destination\nitems by purchase intent (cart=6, order=3, click=1) — straight from the metric.\n"},{"cell_type":"code","metadata":{},"execution_count":null,"outputs":[],"source":"TIME_DELTA_HOURS = 24\nN_TAIL = 30\nTOP_K = 20\nPAIRS_FLUSH = 25_000_000   # we have 30GB now, can push higher\n\ndef _pair_session(aids, ts, types, mode, tdelta):\n    n = aids.shape[0]\n    if n < 2:\n        return (np.empty(0,np.int32),)*2 + (np.empty(0,np.float32),)\n    i, j = np.meshgrid(np.arange(n), np.arange(n), indexing='ij')\n    mask = (i != j) & (np.abs(ts[j] - ts[i]) < tdelta)\n    if not mask.any():\n        return (np.empty(0,np.int32),)*2 + (np.empty(0,np.float32),)\n    ai, aj = i[mask], j[mask]\n    aid_x, aid_y = aids[ai], aids[aj]\n    if mode == 'click2buy':\n        ty = types[aj]\n        w = np.where(ty == 1, 6.0, np.where(ty == 2, 3.0, 1.0)).astype(np.float32)\n    else:\n        w = np.ones(ai.shape[0], dtype=np.float32)\n    return aid_x.astype(np.int32), aid_y.astype(np.int32), w\n\ndef _flush(buf_x, buf_y, buf_w, accum, slack):\n    df = pl.DataFrame({'aid_x': np.concatenate(buf_x),\n                       'aid_y': np.concatenate(buf_y),\n                       'w':     np.concatenate(buf_w)},\n                      schema={'aid_x': pl.Int32,'aid_y': pl.Int32,'w': pl.Float32})\n    df = (df.group_by(['aid_x','aid_y']).agg(pl.col('w').sum())\n            .sort(['aid_x','w'], descending=[False, True])\n            .group_by('aid_x', maintain_order=True).head(slack))\n    if accum is None: return df\n    return (pl.concat([accum, df])\n              .group_by(['aid_x','aid_y']).agg(pl.col('w').sum())\n              .sort(['aid_x','w'], descending=[False, True])\n              .group_by('aid_x', maintain_order=True).head(slack))\n\ndef build_covis(events_pq, mode, out, top_k=TOP_K, n_tail=N_TAIL):\n    t0 = time.time()\n    print(f'>> covis {mode} on {events_pq.name}')\n    lf = pl.scan_parquet(events_pq).select(['session','aid','ts','type'])\n    if mode == 'buy2buy':\n        lf = lf.filter(pl.col('type').is_in([1,2]))\n    df = (lf.sort(['session','ts'], descending=[False, True])\n            .group_by('session', maintain_order=True).head(n_tail)\n            .sort(['session','ts'])\n            .collect())\n    print(f'   tail-loaded events: {df.height:,}')\n    tdelta_unit = 1000 if df['ts'].max() > 10**12 else 1\n    tdelta = TIME_DELTA_HOURS * 3600 * tdelta_unit\n    sess_arr = df['session'].to_numpy()\n    aid_arr  = df['aid'].to_numpy().astype(np.int32)\n    ts_arr   = df['ts'].to_numpy().astype(np.int64)\n    typ_arr  = df['type'].to_numpy().astype(np.int8)\n    change = np.flatnonzero(np.diff(sess_arr)) + 1\n    starts = np.concatenate([[0], change])\n    ends   = np.concatenate([change, [len(sess_arr)]])\n    n_sess = len(starts)\n    print(f'   processing {n_sess:,} sessions...')\n    SLACK = top_k * 4\n    accum = None\n    buf_x, buf_y, buf_w = [], [], []\n    pair_cnt = 0\n    for s in range(n_sess):\n        a, b = starts[s], ends[s]\n        ax, ay, aw = _pair_session(aid_arr[a:b], ts_arr[a:b], typ_arr[a:b], mode, tdelta)\n        if ax.size:\n            buf_x.append(ax); buf_y.append(ay); buf_w.append(aw)\n            pair_cnt += ax.size\n        if pair_cnt >= PAIRS_FLUSH:\n            accum = _flush(buf_x, buf_y, buf_w, accum, SLACK)\n            buf_x.clear(); buf_y.clear(); buf_w.clear(); pair_cnt = 0\n            print(f'   sess {s+1:,}/{n_sess:,}  accum={accum.height:,}  t={time.time()-t0:.1f}s')\n    if buf_x:\n        accum = _flush(buf_x, buf_y, buf_w, accum, SLACK)\n    final = (accum.sort(['aid_x','w'], descending=[False, True])\n                  .group_by('aid_x', maintain_order=True).head(top_k))\n    final.write_parquet(out)\n    print(f'   wrote {out.name}: rows={final.height:,}  anchors={final[\"aid_x\"].n_unique():,}  total={time.time()-t0:.1f}s')\n\n# Use FULL train (includes the validation week — this is fine for the official submission)\nTRAIN_PQ = DATA_FULL/'train.parquet'\nbuild_covis(TRAIN_PQ, 'click2click', ART/'c2c.parquet')\nbuild_covis(TRAIN_PQ, 'click2buy',   ART/'c2b.parquet')\nbuild_covis(TRAIN_PQ, 'buy2buy',     ART/'b2b.parquet')\n"},{"cell_type":"markdown","metadata":{},"source":"## 2. Generate candidates per (session, target)\n\nFor each test session we score each candidate as a weighted sum of:\n- itself (with recency boost) if it appears in the session\n- covis neighbours of recent items (different mix per target)\n- top-20 popularity fallback for sessions with too few candidates\n"},{"cell_type":"code","metadata":{},"execution_count":null,"outputs":[],"source":"TYPE_W = {0: 1.0, 1: 6.0, 2: 3.0}\n\ndef load_covis(path):\n    df = pl.read_parquet(path)\n    out = defaultdict(list)\n    for x, y, w in df.iter_rows():\n        out[x].append((y, w))\n    return out\n\ndef top_pop(events_pq, k=50):\n    return (pl.read_parquet(events_pq, columns=['aid'])\n              .group_by('aid').len()\n              .sort('len', descending=True).head(k)\n              ['aid'].to_list())\n\nprint('loading covis...')\nM_c2c = load_covis(ART/'c2c.parquet')\nM_c2b = load_covis(ART/'c2b.parquet')\nM_b2b = load_covis(ART/'b2b.parquet')\nPOP = top_pop(TRAIN_PQ)\nprint(f'  c2c={len(M_c2c):,}  c2b={len(M_c2b):,}  b2b={len(M_b2b):,}  pop={len(POP)}')\n\ndef gen_candidates(test_pq, label):\n    df = pl.read_parquet(test_pq).sort(['session','ts'])\n    sess = df.group_by('session', maintain_order=True).agg(\n        pl.col('aid').alias('aids'),\n        pl.col('type').alias('types'),\n    )\n    rows_clk, rows_crt, rows_ord = [], [], []\n    t0 = time.time()\n    for i, (sid, aids, types) in enumerate(sess.iter_rows()):\n        own = list(dict.fromkeys(reversed(aids)))\n        types_rev = list(reversed(types))\n        # clicks\n        sc = defaultdict(float)\n        for j, a in enumerate(own):\n            sc[a] += 2.0 / (j + 1)\n            for cand, w in M_c2c.get(a, ()):\n                sc[cand] += w / (j + 1)\n        # popularity backfill\n        if len(sc) < 20:\n            for k, a in enumerate(POP):\n                if a not in sc: sc[a] = -k * 0.001\n        ranked = sorted(sc.items(), key=lambda x: -x[1])[:100]\n        for r, (cand, s) in enumerate(ranked):\n            rows_clk.append((sid, cand, s, r))\n        # carts/orders share recipe\n        sc = defaultdict(float)\n        for j, (a, t) in enumerate(zip(own, types_rev)):\n            mult = TYPE_W.get(t, 1.0)\n            sc[a] += 4.0 * mult / (j + 1)\n            for cand, w in M_c2b.get(a, ()):\n                sc[cand] += w * mult / (j + 1)\n            for cand, w in M_b2b.get(a, ()):\n                sc[cand] += 2.0 * w * mult / (j + 1)\n        if len(sc) < 20:\n            for k, a in enumerate(POP):\n                if a not in sc: sc[a] = -k * 0.001\n        ranked = sorted(sc.items(), key=lambda x: -x[1])[:100]\n        for r, (cand, s) in enumerate(ranked):\n            rows_crt.append((sid, cand, s, r))\n            rows_ord.append((sid, cand, s, r))\n        if (i+1) % 200_000 == 0:\n            print(f'  {i+1:,} sess  t={time.time()-t0:.1f}s')\n\n    schema = {'session': pl.UInt32,'aid': pl.UInt32,'score': pl.Float64,'rank': pl.UInt32}\n    pl.DataFrame(rows_clk, schema=schema, orient='row').write_parquet(ART/f'{label}.clicks.parquet')\n    pl.DataFrame(rows_crt, schema=schema, orient='row').write_parquet(ART/f'{label}.carts.parquet')\n    pl.DataFrame(rows_ord, schema=schema, orient='row').write_parquet(ART/f'{label}.orders.parquet')\n    print(f'  wrote {label}.{{clicks,carts,orders}}.parquet  total={time.time()-t0:.1f}s')\n\nif DATA_RADEK.exists():\n    gen_candidates(DATA_RADEK/'test.parquet', label='val')   # for offline eval\nelse:\n    print('!! DATA_RADEK not attached, skipping offline val candidates')\ngen_candidates(DATA_FULL/'test.parquet',  label='sub')   # official submission\n"},{"cell_type":"markdown","metadata":{},"source":"## 3. Local validation: weighted Recall@20"},{"cell_type":"code","metadata":{},"execution_count":null,"outputs":[],"source":"K = 20\nWEIGHTS = {'clicks': 0.10, 'carts': 0.30, 'orders': 0.60}\nif not DATA_RADEK.exists():\n    print('!! DATA_RADEK not attached — skipping offline recall (will rely on Kaggle LB score after submit)')\nelse:\n    labels = pl.read_parquet(DATA_RADEK/'test_labels.parquet')\n\n    def recall(pred_pq, truth):\n        pred = (pl.read_parquet(pred_pq).sort(['session','rank'])\n                  .group_by('session', maintain_order=True).head(K)\n                  .group_by('session').agg(pl.col('aid').alias('pr')))\n        j = truth.join(pred, on='session', how='left').with_columns(\n            pl.when(pl.col('pr').is_null()).then(pl.lit([], dtype=pl.List(pl.Int64)))\n              .otherwise(pl.col('pr').cast(pl.List(pl.Int64))).alias('pr')\n        )\n        def _r(row):\n            gt, pr = set(row['ground_truth']), set(row['pr'])\n            return len(gt & pr) / min(K, len(gt))\n        return j.with_columns(\n            pl.struct(['ground_truth','pr']).map_elements(_r, return_dtype=pl.Float64).alias('r')\n        )['r'].mean()\n\n    scores = {}\n    for name in ('clicks','carts','orders'):\n        truth = labels.filter(pl.col('type') == name).select(['session','ground_truth'])\n        scores[name] = recall(ART/f'val.{name}.parquet', truth)\n        print(f'{name:8s} R@20 = {scores[name]:.4f}')\n    total = sum(WEIGHTS[k]*v for k,v in scores.items())\n    print(f'WEIGHTED  R@20 = {total:.4f}   (public covis baselines: ~0.557)')\n"},{"cell_type":"markdown","metadata":{},"source":"## 4. Build submission.csv"},{"cell_type":"code","metadata":{},"execution_count":null,"outputs":[],"source":"def build_submission(label, out_csv):\n    rows = []\n    for name in ('clicks','carts','orders'):\n        df = (pl.read_parquet(ART/f'{label}.{name}.parquet')\n                .sort(['session','rank'])\n                .group_by('session', maintain_order=True).head(K))\n        agg = df.group_by('session', maintain_order=True).agg(pl.col('aid').cast(pl.Utf8))\n        for sid, aids in agg.iter_rows():\n            rows.append((f'{sid}_{name}', ' '.join(aids)))\n    pl.DataFrame(rows, schema=['session_type','labels'], orient='row').write_csv(out_csv)\n    print(f'wrote {out_csv} rows={len(rows):,}')\n\nbuild_submission('sub', OUT/'submission.csv')\nimport subprocess\nprint('preview:')\nprint(subprocess.check_output(['head','-5', str(OUT/'submission.csv')]).decode())\nprint('rows:', subprocess.check_output(['wc','-l', str(OUT/'submission.csv')]).decode().strip())\n"},{"cell_type":"markdown","metadata":{},"source":"## 5. Submit\n\nAfter this notebook successfully runs end-to-end, click **Submit to Competition**\nin the right sidebar. It will pick up `submission.csv` from `/kaggle/working/`.\n"}],"metadata":{"kernelspec":{"display_name":"Python 3","language":"python","name":"python3"},"language_info":{"name":"python","version":"3.10"}},"nbformat":4,"nbformat_minor":5}