{"metadata":{"kernelspec":{"language":"python","display_name":"Python 3","name":"python3"},"language_info":{"name":"python","version":"3.12.13","mimetype":"text/x-python","codemirror_mode":{"name":"ipython","version":3},"pygments_lexer":"ipython3","nbconvert_exporter":"python","file_extension":".py"},"kaggle":{"accelerator":"nvidiaRtxPro6000","dataSources":[{"sourceType":"competition","sourceId":129716,"databundleVersionId":16082784,"isSourceIdPinned":false},{"sourceType":"datasetVersion","sourceId":16488582,"datasetId":10545780,"databundleVersionId":17489964},{"sourceType":"datasetVersion","sourceId":16424882,"datasetId":10528332,"databundleVersionId":17421938}],"dockerImageVersionId":31401,"isInternetEnabled":false,"language":"python","sourceType":"notebook","isGpuEnabled":true}},"nbformat_minor":4,"nbformat":4,"cells":[{"cell_type":"code","source":"\"\"\"\nOTTO Recommender System — Co-visitation Matrix + LGBMRanker Pipeline\n\nPhương pháp: Candidate ReRank (2 giai đoạn)\n  1. Candidate Generation: Co-visitation matrices (clicks, cart_order, buy2buy) + history + popular items\n  2. Re-ranking: LGBMRanker (3 models riêng cho clicks/carts/orders)\n\nData: /kaggle/input/datasets/trungnguyen2710t/otto-ptit-filter/\n  - train.parquet   (toàn bộ events để train)\n  - val.parquet     (validation events)\n  - test.parquet    (test events)\n\nMôi trường: Kaggle Notebook (GPU T4)\n\"\"\"\n\n# ══════════════════════════════════════════════════════════════════════════════\n# CELL 1: Imports & Setup\n# ══════════════════════════════════════════════════════════════════════════════\n\nimport os\nimport gc\nimport glob\nimport time\nimport warnings\nfrom pathlib import Path\nfrom collections import defaultdict\n\nimport numpy as np\nimport pandas as pd\nimport polars as pl\nimport lightgbm as lgb\nfrom tqdm import tqdm\n\nwarnings.filterwarnings(\"ignore\")\n\n# GPU engine cho Polars (Kaggle có RAPIDS)\ntry:\n    import cudf_polars\n    print(\"cudf_polars loaded — Polars sẽ chạy trên GPU!\")\n    GPU_ENGINE = True\nexcept ImportError:\n    print(\"cudf_polars không có — Polars chạy trên CPU.\")\n    GPU_ENGINE = False\n\n\ndef _collect(lazy_frame: pl.LazyFrame) -> pl.DataFrame:\n    \"\"\"Collect LazyFrame, dùng GPU engine nếu có.\"\"\"\n    if GPU_ENGINE:\n        try:\n            return lazy_frame.collect(engine=\"gpu\")\n        except Exception:\n            return lazy_frame.collect()\n    return lazy_frame.collect()\n\n\n# ══════════════════════════════════════════════════════════════════════════════\n# CELL 2: Configuration\n# ══════════════════════════════════════════════════════════════════════════════\n\nDATA_PATH = Path(\"/kaggle/input/datasets/trungnguyen2710t/otto-ptit-filter\")\nTRAIN_PARQUET = DATA_PATH / \"train.parquet\"\nVAL_PARQUET   = DATA_PATH / \"val.parquet\"\nTEST_PARQUET  = DATA_PATH / \"test.parquet\"\n\nWORKING_DIR = Path(\"/kaggle/working\")\nWORKING_DIR.mkdir(exist_ok=True)\n\n# --- Hyperparameters ---\nTYPE_LABELS = {\"clicks\": 0, \"carts\": 1, \"orders\": 2}\nTYPE_WEIGHTS_CART_ORDER = {0: 1, 1: 6, 2: 3}  # Trọng số cho cart-order matrix\n\n# Co-visitation matrix params\nCOVISIT_CONFIG = {\n    \"clicks\": {\n        \"top_k\": 20,\n        \"valid_time_hours\": 24,     # Cặp hợp lệ trong 24h\n        \"last_n_events\": 30,        # Chỉ giữ 30 event gần nhất mỗi session\n    },\n    \"cart_order\": {\n        \"top_k\": 15,\n        \"valid_time_hours\": 24,\n        \"last_n_events\": 30,\n    },\n    \"buy2buy\": {\n        \"top_k\": 15,\n        \"valid_time_hours\": 24 * 14,  # 14 ngày\n        \"last_n_events\": 30,\n    },\n}\n\n# LGBMRanker params\nLGBM_PARAMS = {\n    \"objective\": \"lambdarank\",\n    \"metric\": \"map\",\n    \"n_estimators\": 500,\n    \"learning_rate\": 0.05,\n    \"num_leaves\": 31,\n    \"subsample\": 0.8,\n    \"colsample_bytree\": 0.7,\n    \"random_state\": 42,\n    \"n_jobs\": -1,\n    \"verbose\": -1,\n}\n\n# Negative sampling params\nNEG_POS_RATIO = 20        # 20 negative per 1 positive\nPOPULAR_NEG_FRAC = 0.5    # 50% negative = popular (hard negatives)\n\n# Candidate generation\nTOP_POPULAR_CLICKS = 20   # Top popular items for clicks\nTOP_POPULAR_CARTS  = 20   # Top popular items for carts  \nTOP_POPULAR_ORDERS = 20   # Top popular items for orders\nCANDIDATE_CHUNK_SIZE = 50_000  # Sessions per chunk\n\n# Recall metric weights (giống Kaggle)\nRECALL_WEIGHTS = {\"clicks\": 0.10, \"carts\": 0.30, \"orders\": 0.60}\n\nprint(f\"Train : {TRAIN_PARQUET}\")\nprint(f\"Val   : {VAL_PARQUET}\")\nprint(f\"Test  : {TEST_PARQUET}\")\n\n\n# ══════════════════════════════════════════════════════════════════════════════\n# CELL 3: Load & Preprocess Raw Data\n# ══════════════════════════════════════════════════════════════════════════════\n\ndef load_raw_data(parquet_path: Path) -> pl.DataFrame:\n    \"\"\"\n    Load parquet và chuẩn hóa: \n    - type: string → int8 (clicks=0, carts=1, orders=2)\n    - ts: milliseconds → seconds (int32)  \n    - Tự detect xem type là string hay đã là int\n    - Tự detect xem ts đã chia 1000 chưa\n    \"\"\"\n    print(f\"Loading {parquet_path.name}...\")\n    \n    # Load raw trước để kiểm tra schema\n    df = _collect(\n        pl.scan_parquet(str(parquet_path))\n        .select([\"session\", \"aid\", \"ts\", \"type\"])\n    )\n    \n    # --- Xử lý cột type ---\n    # Nếu type là string (\"clicks\", \"carts\", \"orders\") → map sang int\n    # Nếu type đã là int (0, 1, 2) → giữ nguyên\n    if df[\"type\"].dtype == pl.String or df[\"type\"].dtype == pl.Utf8:\n        df = df.with_columns(\n            pl.col(\"type\").replace_strict(TYPE_LABELS).cast(pl.Int8)\n        )\n        print(\"  [type] string → int8 mapping applied\")\n    elif df[\"type\"].dtype not in (pl.Int8, pl.Int16, pl.Int32):\n        df = df.with_columns(pl.col(\"type\").cast(pl.Int8))\n    \n    # --- Xử lý cột ts ---\n    # Nếu ts > 10^12 → đang ở milliseconds, cần chia 1000\n    ts_max = df[\"ts\"].max()\n    if ts_max > 1e12:\n        df = df.with_columns(\n            (pl.col(\"ts\") / 1000).cast(pl.Int32)\n        )\n        print(\"  [ts] milliseconds → seconds conversion applied\")\n    else:\n        df = df.with_columns(pl.col(\"ts\").cast(pl.Int32))\n    \n    print(f\"  → {df.height:,} events | \"\n          f\"{df['session'].n_unique():,} sessions | \"\n          f\"{df['aid'].n_unique():,} unique aids\")\n    return df\n\n\ndef build_ground_truth(val_df: pl.DataFrame) -> pd.DataFrame:\n    \"\"\"\n    Build ground truth từ val.parquet.\n    \n    Dataset structure (otto-ptit-filter):\n    - train.parquet = history events (tuần 1-3)\n    - val.parquet   = future events cần predict (tuần 4)\n    - val_aids.csv  = chỉ chứa danh sách AIDs, KHÔNG phải ground truth\n    \n    Logic: Group val.parquet theo (session, type) → list of unique aids\n    → Đây là ground truth cho Recall@20\n    \n    Returns: pd.DataFrame với columns [session, type, ground_truth]\n    \"\"\"\n    print(\"Building ground truth from val.parquet...\")\n    \n    TYPE_REVERSE = {0: \"clicks\", 1: \"carts\", 2: \"orders\"}\n    \n    gt_pl = (\n        val_df\n        .group_by([\"session\", \"type\"])\n        .agg(pl.col(\"aid\").unique().alias(\"ground_truth\"))\n    )\n    \n    # Convert to pandas cho validation step\n    gt_pd = gt_pl.to_pandas()\n    gt_pd[\"type\"] = gt_pd[\"type\"].map(TYPE_REVERSE)\n    \n    n_types = gt_pd[\"type\"].nunique()\n    n_sessions = gt_pd[\"session\"].nunique()\n    print(f\"  → {len(gt_pd):,} rows | {n_sessions:,} sessions | \"\n          f\"{n_types} types: {gt_pd['type'].unique().tolist()}\")\n    \n    # Stats per type\n    for t in [\"clicks\", \"carts\", \"orders\"]:\n        t_rows = gt_pd[gt_pd[\"type\"] == t]\n        if not t_rows.empty:\n            avg_aids = t_rows[\"ground_truth\"].apply(len).mean()\n            print(f\"    {t:8s}: {len(t_rows):,} sessions | avg {avg_aids:.1f} aids/session\")\n    \n    return gt_pd\n\n\ndef build_train_history_and_labels(\n    df: pl.DataFrame,\n    history_window_days: int = 14,\n) -> tuple:\n    \"\"\"\n    Chia train data thành history (input) và labels (target).\n    \n    Logic:\n    - Cho mỗi session, cutoff = max_ts - 24h\n    - History = events TRƯỚC cutoff (trong history_window_days gần nhất)\n    - Labels  = events SAU cutoff (ground truth)\n    \n    Returns: (history_df, ground_truth_df, full_history_for_features)\n    \"\"\"\n    print(\"\\nBuilding train history & labels...\")\n    \n    # Tính cutoff cho mỗi session\n    session_cutoffs = df.group_by(\"session\").agg(\n        (pl.col(\"ts\").max() - 24 * 3600).alias(\"cutoff_ts\")\n    )\n    \n    last_ts = df[\"ts\"].max()\n    min_ts = last_ts - history_window_days * 24 * 3600\n    \n    # Join cutoff\n    df_with_cutoff = df.join(session_cutoffs, on=\"session\", how=\"left\")\n    \n    # Chia history / labels\n    history_source = df_with_cutoff.filter(\n        (pl.col(\"ts\") < pl.col(\"cutoff_ts\")) & \n        (pl.col(\"ts\") >= min_ts)\n    )\n    labels_source = df_with_cutoff.filter(\n        pl.col(\"ts\") >= pl.col(\"cutoff_ts\")\n    )\n    \n    # Ground truth: group by (session, type) → list of aids\n    ground_truth = (\n        labels_source\n        .group_by([\"session\", \"type\"])\n        .agg(pl.col(\"aid\").alias(\"ground_truth\"))\n    )\n    \n    # History: unique (session, aid) pairs\n    history_df = history_source.select([\"session\", \"aid\"]).unique()\n    \n    # Full history cho feature engineering\n    full_history = history_source.drop(\"cutoff_ts\")\n    \n    print(f\"  History: {history_df.height:,} session-aid pairs | \"\n          f\"{history_df['session'].n_unique():,} sessions\")\n    print(f\"  Labels:  {ground_truth.height:,} ground truth rows\")\n    \n    return history_df, ground_truth, full_history\n\n\n# Load data\ntrain_df = load_raw_data(TRAIN_PARQUET)\nval_df   = load_raw_data(VAL_PARQUET)\n\n# Build ground truth từ val.parquet (KHÔNG dùng val_aids.csv)\nval_gt = build_ground_truth(val_df)\n\n# Kiểm tra overlap giữa train và val sessions\ntrain_sessions = set(train_df[\"session\"].unique().to_list())\nval_sessions = set(val_df[\"session\"].unique().to_list())\noverlap = train_sessions & val_sessions\nprint(f\"\\nSession overlap: {len(overlap):,} sessions appear in BOTH train & val\")\nprint(f\"  Train-only: {len(train_sessions - val_sessions):,}\")\nprint(f\"  Val-only:   {len(val_sessions - train_sessions):,}\")\n\n# Build training history và labels sớm để tránh leakage\ntrain_history, train_gt, train_full_history = build_train_history_and_labels(train_df)\n\ntrain_events_count = train_df.height\n\n# Xóa train_df gốc để giải phóng RAM ngay lập tức\ndel train_df\ngc.collect()\n\n\n# ══════════════════════════════════════════════════════════════════════════════\n# CELL 4: Co-visitation Matrix Generation\n# ══════════════════════════════════════════════════════════════════════════════\n\ndef compute_covisit_matrix(\n    df: pl.DataFrame,\n    matrix_type: str,    # \"clicks\", \"cart_order\", \"buy2buy\"\n    top_k: int = 20,\n    valid_time_hours: int = 24,\n    last_n_events: int = 30,\n    type_weight: dict = None,\n    chunk_size: int = 100_000,  # 100K sessions/chunk (self-join tạo O(N²) pairs)\n) -> pl.DataFrame:\n    \"\"\"\n    Tính ma trận co-visitation trên Polars (CPU).\n    \n    Logic:\n    1. Sort by (session, ts DESC)\n    2. Giữ last_n_events events mỗi session\n    3. Self-join on session → tạo tất cả cặp (aid_x, aid_y)\n    4. Lọc: |ts_x - ts_y| < valid_time VÀ aid_x != aid_y\n    5. Gán trọng số theo loại matrix\n    6. Group by (aid_x, aid_y) → sum(wgt)\n    7. Giữ top_k aid_y cho mỗi aid_x\n    \n    Returns:\n        DataFrame với columns: [aid, candidate_aid, wgt_{type}, rank_{type}]\n    \"\"\"\n    valid_time_sec = valid_time_hours * 3600\n    \n    print(f\"\\n{'='*60}\")\n    print(f\"Computing Co-visitation Matrix: {matrix_type}\")\n    print(f\"  top_k={top_k} | valid_time={valid_time_hours}h | last_n={last_n_events}\")\n    print(f\"{'='*60}\")\n    \n    # Filter cho buy2buy: chỉ carts + orders\n    if matrix_type == \"buy2buy\":\n        source_df = df.filter(pl.col(\"type\").is_in([1, 2]))\n        print(f\"  [buy2buy] Filtered to carts+orders: {source_df.height:,} events\")\n    else:\n        source_df = df\n    \n    # Chia sessions thành chunks để xử lý tuần tự (tránh OOM)\n    all_sessions = source_df[\"session\"].unique().sort().to_list()\n    n_sessions = len(all_sessions)\n    n_chunks = (n_sessions + chunk_size - 1) // chunk_size\n    \n    print(f\"  Processing {n_sessions:,} sessions in {n_chunks} chunks...\")\n    \n    all_results = []\n    \n    for chunk_idx in range(n_chunks):\n        start = chunk_idx * chunk_size\n        end = min(start + chunk_size, n_sessions)\n        chunk_sessions = all_sessions[start:end]\n        \n        t0 = time.time()\n        \n        # Lọc ra events của chunk sessions\n        chunk_df = source_df.filter(pl.col(\"session\").is_in(chunk_sessions))\n        \n        # Sort by (session, ts) DESC để lấy events gần nhất\n        chunk_df = chunk_df.sort([\"session\", \"ts\"], descending=[False, True])\n        \n        # Giữ last_n_events events mỗi session\n        # Dùng row_number thay cho cum_count (tương thích Polars mới hơn)\n        chunk_df = (\n            chunk_df\n            .with_columns(\n                pl.arange(0, pl.len()).over(\"session\").alias(\"event_rank\")\n            )\n            .filter(pl.col(\"event_rank\") < last_n_events)\n            .drop(\"event_rank\")\n        )\n        \n        # Self-join: tạo tất cả cặp trong cùng session\n        pairs_df = chunk_df.join(chunk_df, on=\"session\", suffix=\"_y\")\n        \n        # Lọc cặp hợp lệ\n        pairs_df = pairs_df.filter(\n            ((pl.col(\"ts\") - pl.col(\"ts_y\")).abs() < valid_time_sec) &\n            (pl.col(\"aid\") != pl.col(\"aid_y\"))\n        )\n        \n        # Deduplicate: 1 cặp (session, aid_x, aid_y) chỉ đếm 1 lần\n        pairs_df = pairs_df.unique(subset=[\"session\", \"aid\", \"aid_y\"])\n        \n        # Gán trọng số\n        if matrix_type == \"cart_order\" and type_weight:\n            pairs_df = pairs_df.with_columns(\n                pl.col(\"type_y\").replace_strict(type_weight).cast(pl.Float32).alias(\"wgt\")\n            )\n        elif matrix_type == \"clicks\":\n            # Time-weighted: sự kiện gần hơn có trọng số cao hơn\n            ts_min = pairs_df[\"ts\"].min()\n            ts_max = pairs_df[\"ts\"].max()\n            ts_range = max(ts_max - ts_min, 1)\n            pairs_df = pairs_df.with_columns(\n                (1.0 + 3.0 * (pl.col(\"ts\") - ts_min) / ts_range).cast(pl.Float32).alias(\"wgt\")\n            )\n        else:  # buy2buy\n            pairs_df = pairs_df.with_columns(\n                pl.lit(1.0).cast(pl.Float32).alias(\"wgt\")\n            )\n        \n        # Group by (aid, aid_y) → sum weights\n        result = (\n            pairs_df\n            .group_by([\"aid\", \"aid_y\"])\n            .agg(pl.col(\"wgt\").sum())\n        )\n        \n        all_results.append(result)\n        \n        elapsed = time.time() - t0\n        print(f\"  Chunk {chunk_idx+1}/{n_chunks}: {len(chunk_sessions):,} sessions | \"\n              f\"{result.height:,} pairs | {elapsed:.1f}s\")\n        \n        del chunk_df, pairs_df, result\n        gc.collect()\n    \n    # Merge tất cả chunks\n    print(\"  Merging all chunks...\")\n    merged = pl.concat(all_results)\n    del all_results\n    gc.collect()\n    \n    # Aggregate lại sau merge\n    merged = (\n        merged\n        .group_by([\"aid\", \"aid_y\"])\n        .agg(pl.col(\"wgt\").sum())\n    )\n    \n    # Lấy top_k cho mỗi aid\n    merged = (\n        merged\n        .sort([\"aid\", \"wgt\"], descending=[False, True])\n        .with_columns(\n            pl.arange(0, pl.len()).over(\"aid\").alias(\"rank\")\n        )\n        .filter(pl.col(\"rank\") < top_k)\n    )\n    \n    # Rename columns\n    wgt_col = f\"wgt_{matrix_type}\"\n    rank_col = f\"rank_{matrix_type}\"\n    merged = merged.rename({\n        \"aid_y\": \"candidate_aid\",\n        \"wgt\": wgt_col,\n        \"rank\": rank_col,\n    })\n    \n    # Cast types cho tiết kiệm RAM\n    merged = merged.with_columns(\n        pl.col(rank_col).cast(pl.UInt16),\n    )\n    \n    print(f\"Done: {merged.height:,} total pairs | \"\n          f\"{merged['aid'].n_unique():,} unique source aids\")\n    \n    return merged\n\n\n# --- Tạo 3 ma trận co-visitation ---\nt_start = time.time()\n\ncovisit_clicks = compute_covisit_matrix(\n    train_full_history,\n    matrix_type=\"clicks\",\n    **COVISIT_CONFIG[\"clicks\"],\n)\n\ncovisit_cart_order = compute_covisit_matrix(\n    train_full_history,\n    matrix_type=\"cart_order\",\n    type_weight=TYPE_WEIGHTS_CART_ORDER,\n    **COVISIT_CONFIG[\"cart_order\"],\n)\n\ncovisit_buy2buy = compute_covisit_matrix(\n    train_full_history,\n    matrix_type=\"buy2buy\",\n    **COVISIT_CONFIG[\"buy2buy\"],\n)\n\nprint(f\"\\nTotal co-visitation time: {time.time() - t_start:.0f}s\")\ngc.collect()\n\n\n# ══════════════════════════════════════════════════════════════════════════════\n# CELL 5: Candidate Generation\n# ══════════════════════════════════════════════════════════════════════════════\n\ndef get_popular_items(df: pl.DataFrame, n_days: int = 7) -> pl.DataFrame:\n    \"\"\"Lấy top popular items trong n_days gần nhất.\"\"\"\n    last_ts = df[\"ts\"].max()\n    cutoff_ts = last_ts - n_days * 24 * 3600\n    \n    recent = df.filter(pl.col(\"ts\") >= cutoff_ts)\n    \n    top_clicks = (\n        recent.filter(pl.col(\"type\") == 0)[\"aid\"]\n        .value_counts()\n        .sort(\"count\", descending=True)\n        .head(TOP_POPULAR_CLICKS)[\"aid\"]\n    )\n    top_carts = (\n        recent.filter(pl.col(\"type\") == 1)[\"aid\"]\n        .value_counts()\n        .sort(\"count\", descending=True)\n        .head(TOP_POPULAR_CARTS)[\"aid\"]\n    )\n    top_orders = (\n        recent.filter(pl.col(\"type\") == 2)[\"aid\"]\n        .value_counts()\n        .sort(\"count\", descending=True)\n        .head(TOP_POPULAR_ORDERS)[\"aid\"]\n    )\n    \n    popular = pl.concat([top_clicks, top_carts, top_orders]).unique()\n    popular_df = pl.DataFrame({\"candidate_aid\": popular})\n    print(f\"Popular items: {popular_df.height} unique aids\")\n    return popular_df\n\n\ndef generate_candidates(\n    history_df: pl.DataFrame,    # (session, aid) — unique session-aid pairs\n    df_clicks: pl.DataFrame,     # co-visitation clicks matrix\n    df_buys: pl.DataFrame,       # co-visitation cart_order matrix\n    df_buy2buy: pl.DataFrame,    # co-visitation buy2buy matrix\n    popular_df: pl.DataFrame,    # popular items\n    chunk_size: int = CANDIDATE_CHUNK_SIZE,\n) -> pl.DataFrame:\n    \"\"\"\n    Tổng hợp ứng viên từ 4 nguồn:\n    1. History (items user đã tương tác)\n    2. Popular items\n    3. Co-visitation clicks\n    4. Co-visitation cart-order\n    5. Co-visitation buy2buy\n    \n    Returns: DataFrame (session, candidate_aid) + source features\n    \"\"\"\n    all_sessions = history_df[\"session\"].unique().sort().to_list()\n    n_sessions = len(all_sessions)\n    n_chunks = (n_sessions + chunk_size - 1) // chunk_size\n    \n    print(f\"\\n{'='*60}\")\n    print(f\"Generating Candidates\")\n    print(f\"  {n_sessions:,} sessions in {n_chunks} chunks\")\n    print(f\"{'='*60}\")\n    \n    all_candidates = []\n    \n    for chunk_idx in range(n_chunks):\n        start = chunk_idx * chunk_size\n        end = min(start + chunk_size, n_sessions)\n        chunk_sessions = all_sessions[start:end]\n        \n        t0 = time.time()\n        \n        # History cho chunk này\n        history_chunk = history_df.filter(pl.col(\"session\").is_in(chunk_sessions))\n        \n        # --- Nguồn 1: History candidates ---\n        cand_history = (\n            history_chunk\n            .rename({\"aid\": \"candidate_aid\"})\n            .with_columns(pl.lit(1).cast(pl.UInt8).alias(\"source_history\"))\n        )\n        \n        # --- Nguồn 2: Popular candidates ---\n        sessions_in_chunk = history_chunk.select(\"session\").unique()\n        cand_popular = sessions_in_chunk.join(popular_df, how=\"cross\")\n        \n        # --- Nguồn 3-5: Co-visitation candidates ---\n        cand_clicks_raw = history_chunk.join(df_clicks, on=\"aid\", how=\"inner\")\n        cand_buys_raw   = history_chunk.join(df_buys, on=\"aid\", how=\"inner\")\n        cand_b2b_raw    = history_chunk.join(df_buy2buy, on=\"aid\", how=\"inner\")\n        \n        # --- Union tất cả candidates ---\n        candidates_df = pl.concat([\n            cand_history.select([\"session\", \"candidate_aid\"]),\n            cand_popular.select([\"session\", \"candidate_aid\"]),\n            cand_clicks_raw.select([\"session\", \"candidate_aid\"]),\n            cand_buys_raw.select([\"session\", \"candidate_aid\"]),\n            cand_b2b_raw.select([\"session\", \"candidate_aid\"]),\n        ]).unique(subset=[\"session\", \"candidate_aid\"])\n        \n        # --- Join back source features ---\n        # History flag\n        candidates_df = candidates_df.join(\n            cand_history.select([\"session\", \"candidate_aid\", \"source_history\"]),\n            on=[\"session\", \"candidate_aid\"],\n            how=\"left\",\n        )\n        \n        # Co-visitation features: lấy min rank, max wgt cho mỗi (session, candidate)\n        for cand_raw, name in [\n            (cand_clicks_raw, \"clicks\"),\n            (cand_buys_raw, \"cart_order\"),\n            (cand_b2b_raw, \"buy2buy\"),\n        ]:\n            rank_col = f\"rank_{name}\"\n            wgt_col = f\"wgt_{name}\"\n            \n            if rank_col in cand_raw.columns:\n                agg = (\n                    cand_raw\n                    .group_by([\"session\", \"candidate_aid\"])\n                    .agg([\n                        pl.col(rank_col).min(),\n                        pl.col(wgt_col).max(),\n                    ])\n                )\n                candidates_df = candidates_df.join(\n                    agg, on=[\"session\", \"candidate_aid\"], how=\"left\"\n                )\n        \n        all_candidates.append(candidates_df)\n        \n        elapsed = time.time() - t0\n        n_cand = candidates_df.height\n        avg_per_session = n_cand / max(len(chunk_sessions), 1)\n        print(f\"  Chunk {chunk_idx+1}/{n_chunks}: {n_cand:,} candidates \"\n              f\"({avg_per_session:.0f}/session) | {elapsed:.1f}s\")\n        \n        del (history_chunk, cand_history, cand_popular, \n             cand_clicks_raw, cand_buys_raw, cand_b2b_raw, candidates_df)\n        gc.collect()\n    \n    result = pl.concat(all_candidates)\n    del all_candidates\n    gc.collect()\n    \n    print(f\"Total candidates: {result.height:,} | \"\n          f\"Avg: {result.height / n_sessions:.0f}/session\")\n    \n    return result\n\n\n# --- Generate candidates cho training ---\n\n\n# Sample sessions for training to prevent OOM\nN_TRAIN_SESSIONS = 200_000\nprint(f\"Sampling {N_TRAIN_SESSIONS:,} sessions for training...\")\ntrain_sessions_sampled = train_history[\"session\"].unique().sample(n=N_TRAIN_SESSIONS, seed=42)\ntrain_history = train_history.filter(pl.col(\"session\").is_in(train_sessions_sampled))\ntrain_gt = train_gt.filter(pl.col(\"session\").is_in(train_sessions_sampled))\n\n# Popular items (từ training history)\npopular_items = get_popular_items(train_full_history, n_days=7)\n\n# Generate candidates for training\ntrain_candidates = generate_candidates(\n    train_history, \n    covisit_clicks, \n    covisit_cart_order, \n    covisit_buy2buy,\n    popular_items,\n)\n\ngc.collect()\n\n\n# ══════════════════════════════════════════════════════════════════════════════\n# CELL 6: Feature Engineering\n# ══════════════════════════════════════════════════════════════════════════════\n\ndef create_item_features(df: pl.DataFrame, suffix: str = \"_all\") -> pl.DataFrame:\n    \"\"\"Tạo item-level features (count + ratio) cho một cửa sổ thời gian.\"\"\"\n    item_feats = (\n        df\n        .group_by(\"aid\")\n        .agg([\n            pl.count().alias(f\"item_total_cnt{suffix}\"),\n            pl.col(\"type\").filter(pl.col(\"type\") == 0).count().alias(f\"item_click_cnt{suffix}\"),\n            pl.col(\"type\").filter(pl.col(\"type\") == 1).count().alias(f\"item_cart_cnt{suffix}\"),\n            pl.col(\"type\").filter(pl.col(\"type\") == 2).count().alias(f\"item_order_cnt{suffix}\"),\n            pl.col(\"session\").n_unique().alias(f\"item_unique_sessions{suffix}\"),\n        ])\n        .rename({\"aid\": \"candidate_aid\"})\n    )\n    \n    # Conversion ratios (smoothing +10)\n    item_feats = item_feats.with_columns([\n        (pl.col(f\"item_order_cnt{suffix}\") / (pl.col(f\"item_click_cnt{suffix}\") + 10))\n            .alias(f\"item_buy_ratio{suffix}\"),\n        (pl.col(f\"item_cart_cnt{suffix}\") / (pl.col(f\"item_click_cnt{suffix}\") + 10))\n            .alias(f\"item_cart_ratio{suffix}\"),\n        (pl.col(f\"item_order_cnt{suffix}\") / (pl.col(f\"item_cart_cnt{suffix}\") + 10))\n            .alias(f\"item_order_per_cart{suffix}\"),\n    ])\n    \n    return item_feats\n\n\ndef create_session_features(df: pl.DataFrame) -> tuple:\n    \"\"\"\n    Tạo session-level features.\n    Returns: (session_feats, interaction_feats, last_item_info)\n    \"\"\"\n    # Session level\n    session_feats = (\n        df\n        .group_by(\"session\")\n        .agg([\n            pl.count().alias(\"session_length\"),\n            pl.col(\"aid\").n_unique().alias(\"session_unique_aids\"),\n            pl.col(\"ts\").max().alias(\"session_end_ts\"),\n            (pl.col(\"ts\").max() - pl.col(\"ts\").min()).alias(\"session_duration\"),\n            # Type distribution\n            pl.col(\"type\").filter(pl.col(\"type\") == 0).count().alias(\"session_click_cnt\"),\n            pl.col(\"type\").filter(pl.col(\"type\") == 1).count().alias(\"session_cart_cnt\"),\n            pl.col(\"type\").filter(pl.col(\"type\") == 2).count().alias(\"session_order_cnt\"),\n        ])\n    )\n    \n    # Click ratio trong session\n    session_feats = session_feats.with_columns([\n        (pl.col(\"session_click_cnt\") / (pl.col(\"session_length\") + 1))\n            .alias(\"session_click_ratio\"),\n        (pl.col(\"session_cart_cnt\") / (pl.col(\"session_length\") + 1))\n            .alias(\"session_cart_ratio\"),\n    ])\n    \n    # Interaction level (per session-aid pair)\n    interaction_feats = (\n        df\n        .group_by([\"session\", \"aid\"])\n        .agg([\n            pl.count().alias(\"num_repetitions\"),\n            pl.col(\"ts\").max().alias(\"last_item_ts\"),\n            pl.col(\"ts\").min().alias(\"first_item_ts\"),\n        ])\n        .rename({\"aid\": \"candidate_aid\"})\n    )\n    \n    # Last item in session\n    last_items = (\n        df\n        .sort(\"ts\")\n        .group_by(\"session\", maintain_order=True)\n        .last()\n        .select([\"session\", \"aid\"])\n        .rename({\"aid\": \"last_aid\"})\n    )\n    \n    return session_feats, interaction_feats, last_items\n\n\ndef add_all_features(\n    candidates_df: pl.DataFrame,\n    feature_source_df: pl.DataFrame,\n    global_item_feats_all: pl.DataFrame,\n    global_item_feats_7d: pl.DataFrame,\n) -> pl.DataFrame:\n    \"\"\"\n    Thêm tất cả features cho candidates DataFrame.\n    \"\"\"\n    print(\"\\nAdding features...\")\n    t0 = time.time()\n    \n    df = candidates_df.clone()\n    \n    # --- Item features ---\n    df = df.join(global_item_feats_all, on=\"candidate_aid\", how=\"left\")\n    df = df.join(global_item_feats_7d, on=\"candidate_aid\", how=\"left\")\n    \n    # --- Session features ---\n    # Lọc feature_source_df để chỉ tính session features cho các session trong candidates\n    active_sessions = df[\"session\"].unique()\n    session_history = feature_source_df.filter(pl.col(\"session\").is_in(active_sessions))\n    \n    session_feats, interaction_feats, last_items = create_session_features(session_history)\n    \n    df = df.join(session_feats, on=\"session\", how=\"left\")\n    df = df.join(interaction_feats, on=[\"session\", \"candidate_aid\"], how=\"left\")\n    df = df.join(last_items, on=\"session\", how=\"left\")\n    \n    # --- Derived features ---\n    \n    # Fill null ranks with 999\n    for c in [\"rank_clicks\", \"rank_cart_order\", \"rank_buy2buy\"]:\n        if c not in df.columns:\n            df = df.with_columns(pl.lit(999).cast(pl.UInt16).alias(c))\n        else:\n            df = df.with_columns(pl.col(c).fill_null(999))\n    \n    # Fill null weights with 0\n    for c in [\"wgt_clicks\", \"wgt_cart_order\", \"wgt_buy2buy\"]:\n        if c not in df.columns:\n            df = df.with_columns(pl.lit(0.0).cast(pl.Float32).alias(c))\n        else:\n            df = df.with_columns(pl.col(c).fill_null(0.0))\n    \n    df = df.with_columns([\n        # --- Trend features ---\n        (pl.col(\"item_click_cnt_7d\").fill_null(0) / (pl.col(\"item_click_cnt_all\").fill_null(0) + 10))\n            .alias(\"click_trend_7d\"),\n        (pl.col(\"item_order_cnt_7d\").fill_null(0) / (pl.col(\"item_order_cnt_all\").fill_null(0) + 10))\n            .alias(\"order_trend_7d\"),\n        (pl.col(\"item_buy_ratio_7d\").fill_null(0) - pl.col(\"item_buy_ratio_all\").fill_null(0))\n            .alias(\"conversion_trend_diff\"),\n        \n        # --- Cross-source rank features ---\n        (pl.col(\"rank_clicks\") - pl.col(\"rank_buy2buy\")).cast(pl.Int32)\n            .alias(\"rank_diff_click_b2b\"),\n        (pl.col(\"rank_cart_order\") - pl.col(\"rank_buy2buy\")).cast(pl.Int32)\n            .alias(\"rank_diff_buy_b2b\"),\n        \n        # Combined weight\n        (pl.col(\"wgt_buy2buy\") * 2.0 + pl.col(\"wgt_cart_order\") * 1.0)\n            .alias(\"combined_buy_weight\"),\n        \n        # --- Recency features ---\n        (pl.col(\"session_end_ts\") - pl.col(\"last_item_ts\")).fill_null(7 * 24 * 3600)\n            .alias(\"recency_sec\"),\n        \n        # --- Flags ---\n        (pl.col(\"candidate_aid\") == pl.col(\"last_aid\")).cast(pl.Int8).fill_null(0)\n            .alias(\"is_last_viewed\"),\n        (pl.col(\"num_repetitions\").fill_null(0) > 1).cast(pl.Int8)\n            .alias(\"is_repeated\"),\n        \n        # Source history flag\n        pl.col(\"source_history\").fill_null(0),\n    ])\n    \n    # Log recency\n    df = df.with_columns(\n        pl.col(\"recency_sec\").cast(pl.Float64).log1p().alias(\"log_recency\")\n    )\n    \n    # Decayed weights\n    df = df.with_columns([\n        (pl.col(\"wgt_buy2buy\") / (pl.col(\"log_recency\") + 1)).alias(\"wgt_b2b_decayed\"),\n        (pl.col(\"wgt_clicks\") / (pl.col(\"log_recency\") + 1)).alias(\"wgt_clicks_decayed\"),\n    ])\n    \n    # Sorted rank features: min_rank_1, min_rank_2, n_sources (Optimized using horizontal column-wise ops)\n    a = pl.col(\"rank_clicks\").cast(pl.Int32)\n    b = pl.col(\"rank_cart_order\").cast(pl.Int32)\n    c = pl.col(\"rank_buy2buy\").cast(pl.Int32)\n    \n    min_val = pl.min_horizontal([a, b, c])\n    max_val = pl.max_horizontal([a, b, c])\n    \n    df = df.with_columns([\n        min_val.cast(pl.UInt16).alias(\"min_rank_1\"),\n        (a + b + c - min_val - max_val).cast(pl.UInt16).alias(\"min_rank_2\"),\n        ((a < 999).cast(pl.Int8) + (b < 999).cast(pl.Int8) + (c < 999).cast(pl.Int8)).alias(\"n_sources_present\")\n    ])\n\n    \n    # Drop helper columns\n    cols_to_drop = [\"session_end_ts\", \"last_item_ts\", \"first_item_ts\", \"last_aid\"]\n    df = df.drop([c for c in cols_to_drop if c in df.columns])\n    \n    # Fill remaining nulls\n    df = df.fill_null(0)\n    \n    elapsed = time.time() - t0\n    print(f\"Features added: {len(df.columns)} columns | {elapsed:.1f}s\")\n    \n    return df\n\n\n# --- Precompute global item features ---\nprint(\"\\nPrecomputing global item features on train_full_history...\")\nt_item_start = time.time()\nlast_ts = train_full_history[\"ts\"].max()\nglobal_item_feats_all = create_item_features(train_full_history, suffix=\"_all\")\nrecent_7d = train_full_history.filter(pl.col(\"ts\") >= last_ts - 7 * 24 * 3600)\nglobal_item_feats_7d = create_item_features(recent_7d, suffix=\"_7d\")\nprint(f\"  Done precomputing global item features: {time.time() - t_item_start:.1f}s\")\ngc.collect()\n\n\n# --- Thêm features cho training candidates ---\ntrain_candidates_featured = add_all_features(\n    train_candidates, \n    train_full_history,\n    global_item_feats_all,\n    global_item_feats_7d,\n)\n\ngc.collect()\nprint(f\"\\nFinal train candidates shape: {train_candidates_featured.shape}\")\nprint(f\"Columns: {train_candidates_featured.columns}\")\n\n\n# ══════════════════════════════════════════════════════════════════════════════\n# CELL 7: Training Set Construction\n# ══════════════════════════════════════════════════════════════════════════════\n\ndef create_labeled_training_set(\n    candidates_df: pl.DataFrame,\n    ground_truth_df: pl.DataFrame,\n    prediction_type: str,     # \"clicks\", \"carts\", \"orders\"\n    neg_pos_ratio: int = NEG_POS_RATIO,\n    popular_neg_frac: float = POPULAR_NEG_FRAC,\n    seed: int = 42,\n) -> pl.DataFrame:\n    \"\"\"\n    Tạo training set có label cho một loại prediction.\n    \n    Logic:\n    1. Join candidates với ground truth → label=1 nếu khớp, else 0\n    2. Lấy toàn bộ positives\n    3. Sample negatives: 50% popular (hard) + 50% random (easy)\n    \"\"\"\n    type_code = TYPE_LABELS[prediction_type]\n    \n    print(f\"\\n--- Creating training set for '{prediction_type}' (type={type_code}) ---\")\n    \n    # Lấy ground truth cho type này\n    # ground_truth_df (Polars) có cột type là Int8, so sánh bằng int\n    type_gt = ground_truth_df.filter(pl.col(\"type\") == type_code)\n    \n    if type_gt.height == 0:\n        print(f\"No ground truth found for type={type_code}. Check type column values.\")\n        print(f\"     Available types: {ground_truth_df['type'].unique().to_list()}\")\n        # Fallback: return empty labeled candidates\n        return candidates_df.head(0).with_columns(pl.lit(0).cast(pl.UInt8).alias(\"label\"))\n    \n    # Explode ground_truth list → individual (session, candidate_aid)\n    type_labels = (\n        type_gt\n        .explode(\"ground_truth\")\n        .rename({\"ground_truth\": \"candidate_aid\"})\n        .with_columns(pl.lit(1).cast(pl.UInt8).alias(\"label\"))\n        .select([\"session\", \"candidate_aid\", \"label\"])\n        .unique()\n    )\n    \n    # Join labels\n    labeled_df = candidates_df.join(\n        type_labels,\n        on=[\"session\", \"candidate_aid\"],\n        how=\"left\",\n    ).with_columns(pl.col(\"label\").fill_null(0))\n    \n    # Tách positive / negative\n    positives = labeled_df.filter(pl.col(\"label\") == 1)\n    negatives = labeled_df.filter(pl.col(\"label\") == 0)\n    \n    n_pos = positives.height\n    n_neg_target = min(n_pos * neg_pos_ratio, negatives.height)\n    n_popular_neg = int(n_neg_target * popular_neg_frac)\n    n_random_neg = n_neg_target - n_popular_neg\n    \n    print(f\"  Positives: {n_pos:,}\")\n    print(f\"  Negatives target: {n_neg_target:,} ({n_popular_neg:,} popular + {n_random_neg:,} random)\")\n    \n    # Popular negatives: lấy những negatives có item_total_cnt cao\n    if \"item_total_cnt_all\" in negatives.columns:\n        popular_neg = (\n            negatives\n            .sort(\"item_total_cnt_all\", descending=True)\n            .head(n_popular_neg)\n        )\n    else:\n        popular_neg = negatives.sample(n=min(n_popular_neg, negatives.height), seed=seed)\n    \n    # Random negatives: từ phần còn lại\n    remaining = negatives.join(\n        popular_neg.select([\"session\", \"candidate_aid\"]),\n        on=[\"session\", \"candidate_aid\"],\n        how=\"anti\",\n    )\n    random_neg = remaining.sample(\n        n=min(n_random_neg, remaining.height),\n        seed=seed,\n    )\n    \n    # Concat\n    final_cols = positives.columns\n    final_df = pl.concat([\n        positives,\n        popular_neg.select(final_cols),\n        random_neg.select(final_cols),\n    ])\n    \n    print(f\"  Final training set: {final_df.height:,} rows \"\n          f\"(pos/neg ratio = 1:{(final_df.height - n_pos) / max(n_pos, 1):.1f})\")\n    \n    return final_df\n\n\n# Tạo training sets cho 3 types\ntrain_sets = {}\nfor pred_type in [\"clicks\", \"carts\", \"orders\"]:\n    train_sets[pred_type] = create_labeled_training_set(\n        train_candidates_featured,\n        train_gt,\n        prediction_type=pred_type,\n    )\n\ngc.collect()\n\n\n# ══════════════════════════════════════════════════════════════════════════════\n# CELL 8: LGBMRanker Training\n# ══════════════════════════════════════════════════════════════════════════════\n\ndef get_feature_columns(df: pl.DataFrame) -> list:\n    \"\"\"Lấy danh sách feature columns (loại bỏ session, candidate_aid, label).\"\"\"\n    ignore = {\"session\", \"candidate_aid\", \"label\"}\n    return [c for c in df.columns if c not in ignore]\n\n\ndef train_lgbm_ranker(\n    df: pl.DataFrame,\n    model_type: str,\n    params: dict = LGBM_PARAMS,\n) -> tuple:\n    \"\"\"\n    Train LGBMRanker cho một loại prediction.\n    \n    Returns: (model, feature_names)\n    \"\"\"\n    print(f\"\\n{'='*60}\")\n    print(f\"Training LGBMRanker for '{model_type}'\")\n    print(f\"{'='*60}\")\n    \n    # Sort by session (QUAN TRỌNG cho LGBMRanker)\n    df = df.sort(\"session\")\n    \n    feature_cols = get_feature_columns(df)\n    print(f\"  Features: {len(feature_cols)}\")\n    \n    X = df.select(feature_cols).to_numpy()\n    y = df.select(\"label\").to_numpy().ravel()\n    \n    # Groups: số candidates mỗi session\n    groups = df.group_by(\"session\", maintain_order=True).len()[\"len\"].to_numpy()\n    \n    print(f\"  X shape: {X.shape}\")\n    print(f\"  Positives: {y.sum():,.0f} | Negatives: {(1-y).sum():,.0f}\")\n    print(f\"  Groups: {len(groups):,} sessions\")\n    \n    # Train\n    model = lgb.LGBMRanker(**params)\n    \n    model.fit(\n        X, y,\n        group=groups,\n        eval_set=[(X, y)],\n        eval_group=[groups],\n        callbacks=[\n            lgb.early_stopping(50, verbose=False),\n            lgb.log_evaluation(100),\n        ],\n    )\n    \n    # Feature importance\n    imp = sorted(\n        zip(feature_cols, model.feature_importances_),\n        key=lambda x: x[1],\n        reverse=True,\n    )\n    print(f\"\\n  Top 10 features:\")\n    for name, gain in imp[:10]:\n        print(f\"    {name:35s} → {gain:,.0f}\")\n    \n    return model, feature_cols\n\n\n# Train 3 models\nmodels = {}\nfeature_names = {}\n\nfor pred_type in [\"clicks\", \"carts\", \"orders\"]:\n    models[pred_type], feature_names[pred_type] = train_lgbm_ranker(\n        train_sets[pred_type],\n        model_type=pred_type,\n    )\n\ngc.collect()\n\n\n# ══════════════════════════════════════════════════════════════════════════════\n# CELL 9: Validation (Local CV)\n# ══════════════════════════════════════════════════════════════════════════════\n\ndef run_validation(\n    val_df: pl.DataFrame,\n    val_gt_pd: pd.DataFrame,\n    train_df_for_features: pl.DataFrame,\n    models: dict,\n    feature_names: dict,\n    global_item_feats_all: pl.DataFrame,\n    global_item_feats_7d: pl.DataFrame,\n    chunk_size: int = 50_000,\n) -> dict:\n    \"\"\"\n    Chạy full pipeline trên validation set (theo chunks) và tính Recall@20.\n    \n    Data flow:\n    - train.parquet = history (events TRƯỚC cutoff)  → dùng làm input\n    - val.parquet   = future  (events SAU cutoff)    → dùng làm ground truth\n    - History cho validation = train_df events của sessions xuất hiện trong val_df\n    \"\"\"\n    print(f\"\\n{'='*60}\")\n    print(\"Running Validation (Chunked)\")\n    print(f\"{'='*60}\")\n    \n    # 1. Build val history\n    # History = events từ TRAIN cho các sessions có trong VAL\n    # (vì val.parquet chứa future events = ground truth)\n    val_session_ids = val_df[\"session\"].unique().sort().to_list()\n    val_history_source = train_df_for_features.filter(\n        pl.col(\"session\").is_in(val_session_ids)\n    )\n    \n    # Nếu không có overlap → val sessions không có trong train\n    # → dùng val_df events trước cutoff làm history, events sau cutoff làm ground truth\n    if val_history_source.height == 0:\n        print(\" No session overlap between train & val!\")\n        print(\"  → Falling back: splitting val.parquet internally (first 80% = history / last 20% = ground truth)\")\n        # Fallback: chia val.parquet theo thời gian\n        val_cutoff = val_df.group_by(\"session\").agg(\n            (pl.col(\"ts\").min() + (pl.col(\"ts\").max() - pl.col(\"ts\").min()) * 0.8)\n                .cast(pl.Int32).alias(\"cutoff\")\n        )\n        val_with_cutoff = val_df.join(val_cutoff, on=\"session\")\n        val_history_source = val_with_cutoff.filter(pl.col(\"ts\") < pl.col(\"cutoff\")).drop(\"cutoff\")\n        val_labels_source = val_with_cutoff.filter(pl.col(\"ts\") >= pl.col(\"cutoff\")).drop(\"cutoff\")\n        \n        # Build lại ground truth chỉ từ phần labels sau cutoff\n        val_gt_pd = build_ground_truth(val_labels_source)\n    \n    val_history = val_history_source.select([\"session\", \"aid\"]).unique()\n    print(f\"Val history: {val_history.height:,} session-aid pairs | \"\n          f\"{val_history['session'].n_unique():,} sessions\")\n    \n    n_sessions = len(val_session_ids)\n    n_chunks = (n_sessions + chunk_size - 1) // chunk_size\n    print(f\"Processing validation in {n_chunks} chunks of {chunk_size:,} sessions...\")\n    \n    all_top_20 = {\n        \"clicks\": [],\n        \"carts\": [],\n        \"orders\": []\n    }\n    \n    for chunk_idx in range(n_chunks):\n        start = chunk_idx * chunk_size\n        end = min(start + chunk_size, n_sessions)\n        chunk_sessions = val_session_ids[start:end]\n        \n        t0 = time.time()\n        \n        # Lọc history cho chunk này\n        chunk_hist = val_history.filter(pl.col(\"session\").is_in(chunk_sessions))\n        chunk_hist_src = val_history_source.filter(pl.col(\"session\").is_in(chunk_sessions))\n        \n        # 2. Generate candidates\n        chunk_candidates = generate_candidates(\n            chunk_hist,\n            covisit_clicks,\n            covisit_cart_order,\n            covisit_buy2buy,\n            popular_items,\n            chunk_size=chunk_size,\n        )\n        \n        # 3. Add features\n        chunk_featured = add_all_features(\n            chunk_candidates, \n            chunk_hist_src,\n            global_item_feats_all,\n            global_item_feats_7d,\n        )\n        \n        # 4. Predict scores\n        chunk_preds = chunk_featured.select([\"session\", \"candidate_aid\"]).clone()\n        \n        for pred_type in [\"clicks\", \"carts\", \"orders\"]:\n            feats = feature_names[pred_type]\n            # Đảm bảo tất cả feature columns tồn tại\n            for f in feats:\n                if f not in chunk_featured.columns:\n                    chunk_featured = chunk_featured.with_columns(pl.lit(0.0).alias(f))\n            \n            X_val = chunk_featured.select(feats).to_numpy()\n            scores = models[pred_type].predict(X_val)\n            \n            preds_type_df = chunk_featured.select([\"session\", \"candidate_aid\"]).with_columns(\n                pl.Series(\"score\", scores)\n            )\n            \n            top_20 = (\n                preds_type_df\n                .sort(\"score\", descending=True)\n                .group_by(\"session\", maintain_order=False)\n                .head(20)\n                .group_by(\"session\")\n                .agg(pl.col(\"candidate_aid\").alias(\"predicted_aids\"))\n            )\n            all_top_20[pred_type].append(top_20)\n            \n        elapsed = time.time() - t0\n        print(f\"  Chunk {chunk_idx+1}/{n_chunks}: {len(chunk_sessions):,} sessions completed in {elapsed:.1f}s\")\n        \n        del chunk_hist, chunk_hist_src, chunk_candidates, chunk_featured, chunk_preds\n        gc.collect()\n    \n    # 5. Tính Recall@20\n    recalls = {}\n    for pred_type in [\"clicks\", \"carts\", \"orders\"]:\n        print(f\"\\nCalculating Recall@20 for '{pred_type}'...\")\n        # Combine all chunk predictions\n        top_20_all = pl.concat(all_top_20[pred_type])\n        top_20_all_pd = top_20_all.to_pandas()\n        \n        # Merge with ground truth\n        # val_gt_pd[\"type\"] is string, pred_type is also string → match directly\n        gt_type = val_gt_pd[val_gt_pd[\"type\"] == pred_type].copy()\n        \n        if gt_type.empty:\n            # Try matching with integer type code\n            gt_type = val_gt_pd[val_gt_pd[\"type\"] == str(TYPE_LABELS[pred_type])].copy()\n        \n        if gt_type.empty:\n            print(f\" No ground truth for {pred_type}\")\n            print(f\"     Available types: {val_gt_pd['type'].unique().tolist()}\")\n            recalls[pred_type] = 0.0\n            continue\n        \n        merged = gt_type.merge(top_20_all_pd, on=\"session\", how=\"left\")\n        \n        predicted = [list(x) if hasattr(x, '__iter__') else [] for x in merged[\"predicted_aids\"]]\n        ground_truth = [list(x) if hasattr(x, '__iter__') else [] for x in merged[\"ground_truth\"]]\n        \n        # Calculate hits\n        hits = [\n            len(set(gt).intersection(set(pred)))\n            for gt, pred in zip(ground_truth, predicted)\n        ]\n        gt_counts = [min(len(gt), 20) for gt in ground_truth]\n        \n        recall = sum(hits) / max(sum(gt_counts), 1)\n        recalls[pred_type] = recall\n        print(f\"  {pred_type:8s} Recall@20 = {recall:.5f}\")\n    \n    # Weighted score\n    weighted_score = sum(\n        RECALL_WEIGHTS[t] * recalls[t] for t in recalls\n    )\n    print(f\"\\n  Weighted Score = {weighted_score:.5f}\")\n    print(f\"    (0.10×clicks + 0.30×carts + 0.60×orders)\")\n    \n    return recalls\n\n\n# Run validation\nval_recalls = run_validation(\n    val_df, val_gt,\n    train_full_history,  # Dùng train_full_history làm features (train_df đã được xóa để giải phóng RAM)\n    models, feature_names,\n    global_item_feats_all,\n    global_item_feats_7d,\n)\n\n\n# ══════════════════════════════════════════════════════════════════════════════\n# CELL 9.5: Serialization & Standalone Batch Inference Utilities\n# ══════════════════════════════════════════════════════════════════════════════\n\ndef save_pipeline_artifacts(\n    covisit_clicks: pl.DataFrame,\n    covisit_cart_order: pl.DataFrame,\n    covisit_buy2buy: pl.DataFrame,\n    popular_items: pl.DataFrame,\n    global_item_feats_all: pl.DataFrame,\n    global_item_feats_7d: pl.DataFrame,\n    models: dict,\n    feature_names: dict,\n    artifacts_dir: str = \"/kaggle/working/artifacts\"\n):\n    \"\"\"\n    Lưu tất cả co-visitation matrices, features, LightGBM models, và metadata xuống đĩa.\n    \"\"\"\n    dir_path = Path(artifacts_dir)\n    dir_path.mkdir(parents=True, exist_ok=True)\n    \n    print(f\"\\nSaving pipeline artifacts to {dir_path}...\")\n    t0 = time.time()\n    \n    # Lưu Polars DataFrames dưới dạng parquet (nén tốt, đọc ghi cực nhanh)\n    covisit_clicks.write_parquet(dir_path / \"covisit_clicks.parquet\")\n    covisit_cart_order.write_parquet(dir_path / \"covisit_cart_order.parquet\")\n    covisit_buy2buy.write_parquet(dir_path / \"covisit_buy2buy.parquet\")\n    popular_items.write_parquet(dir_path / \"popular_items.parquet\")\n    global_item_feats_all.write_parquet(dir_path / \"global_item_feats_all.parquet\")\n    global_item_feats_7d.write_parquet(dir_path / \"global_item_feats_7d.parquet\")\n    \n    # Lưu các model LGBMRanker bằng pickle\n    import pickle\n    with open(dir_path / \"models.pkl\", \"wb\") as f:\n        pickle.dump(models, f)\n        \n    # Lưu feature names bằng json\n    import json\n    with open(dir_path / \"feature_names.json\", \"w\") as f:\n        json.dump(feature_names, f, indent=4)\n        \n    print(f\"Artifacts saved successfully in {time.time() - t0:.1f}s!\")\n\n\ndef load_pipeline_artifacts(\n    artifacts_dir: str = \"/kaggle/working/artifacts\"\n) -> tuple:\n    \"\"\"\n    Load lại toàn bộ artifacts đã được lưu từ đĩa.\n    \"\"\"\n    dir_path = Path(artifacts_dir)\n    if not dir_path.exists():\n        raise FileNotFoundError(f\"Artifacts directory {dir_path} not found.\")\n        \n    print(f\"\\nLoading pipeline artifacts from {dir_path}...\")\n    t0 = time.time()\n    \n    covisit_clicks = pl.read_parquet(dir_path / \"covisit_clicks.parquet\")\n    covisit_cart_order = pl.read_parquet(dir_path / \"covisit_cart_order.parquet\")\n    covisit_buy2buy = pl.read_parquet(dir_path / \"covisit_buy2buy.parquet\")\n    popular_items = pl.read_parquet(dir_path / \"popular_items.parquet\")\n    global_item_feats_all = pl.read_parquet(dir_path / \"global_item_feats_all.parquet\")\n    global_item_feats_7d = pl.read_parquet(dir_path / \"global_item_feats_7d.parquet\")\n    \n    import pickle\n    with open(dir_path / \"models.pkl\", \"rb\") as f:\n        models = pickle.load(f)\n        \n    import json\n    with open(dir_path / \"feature_names.json\", \"r\") as f:\n        feature_names = json.load(f)\n        \n    print(f\"Artifacts loaded successfully in {time.time() - t0:.1f}s!\")\n    return (\n        covisit_clicks,\n        covisit_cart_order,\n        covisit_buy2buy,\n        popular_items,\n        global_item_feats_all,\n        global_item_feats_7d,\n        models,\n        feature_names\n    )\n\n\ndef recommend_for_batch(\n    session_events: pl.DataFrame,\n    covisit_clicks: pl.DataFrame,\n    covisit_cart_order: pl.DataFrame,\n    covisit_buy2buy: pl.DataFrame,\n    popular_items: pl.DataFrame,\n    global_item_feats_all: pl.DataFrame,\n    global_item_feats_7d: pl.DataFrame,\n    models: dict,\n    feature_names: dict,\n) -> dict:\n    \"\"\"\n    Nhận vào DataFrame chứa lịch sử tương tác của một lô sessions (session, aid, ts, type).\n    Trả về gợi ý top 20 (clicks, carts, orders) cho từng session dưới dạng python dict:\n    {\n        session_id: {\n            \"clicks\": [aid1, aid2, ...],\n            \"carts\": [aid1, aid2, ...],\n            \"orders\": [aid1, aid2, ...]\n        },\n        ...\n    }\n    \"\"\"\n    sessions = session_events[\"session\"].unique().to_list()\n    results = {sess: {\"clicks\": [], \"carts\": [], \"orders\": []} for sess in sessions}\n    \n    # 1. Trích xuất history\n    history = session_events.select([\"session\", \"aid\"]).unique()\n    \n    # 2. Tạo ứng viên\n    candidates = generate_candidates(\n        history,\n        covisit_clicks,\n        covisit_cart_order,\n        covisit_buy2buy,\n        popular_items,\n        chunk_size=max(len(sessions), 1)\n    )\n    \n    if candidates.height == 0:\n        return results\n        \n    # 3. Thêm features\n    featured = add_all_features(\n        candidates,\n        session_events,\n        global_item_feats_all,\n        global_item_feats_7d,\n    )\n    \n    if featured.height == 0:\n        return results\n        \n    # 4. Dự đoán điểm số và lọc top 20 cho mỗi mục tiêu\n    for pred_type in [\"clicks\", \"carts\", \"orders\"]:\n        feats = feature_names[pred_type]\n        # Đảm bảo các feature tồn tại\n        for f in feats:\n            if f not in featured.columns:\n                featured = featured.with_columns(pl.lit(0.0).alias(f))\n                \n        X = featured.select(feats).to_numpy()\n        scores = models[pred_type].predict(X)\n        \n        preds_df = featured.select([\"session\", \"candidate_aid\"]).with_columns(\n            pl.Series(\"score\", scores)\n        )\n        \n        top_20 = (\n            preds_df\n            .sort(\"score\", descending=True)\n            .group_by(\"session\", maintain_order=False)\n            .head(20)\n            .group_by(\"session\")\n            .agg(pl.col(\"candidate_aid\").alias(\"aids\"))\n        )\n        \n        for row in top_20.iter_rows():\n            sess_id, aids = row\n            if sess_id in results:\n                results[sess_id][pred_type] = list(aids)\n                \n    return results\n\n\n# ══════════════════════════════════════════════════════════════════════════════\n# CELL 10: Test Inference & Submission\n# ══════════════════════════════════════════════════════════════════════════════\n\ndef create_submission(\n    test_df: pl.DataFrame,\n    train_df: pl.DataFrame,\n    models: dict,\n    feature_names: dict,\n    global_item_feats_all: pl.DataFrame,\n    global_item_feats_7d: pl.DataFrame,\n    output_path: str = \"/kaggle/working/submission.csv\",\n    chunk_size: int = 50_000,\n) -> pl.DataFrame:\n    \"\"\"\n    Chạy full inference trên test set và tạo file submission theo chunks.\n    \"\"\"\n    print(f\"\\n{'='*60}\")\n    print(\"Creating Submission (Chunked)\")\n    print(f\"{'='*60}\")\n    \n    # 1. History\n    test_history = test_df.select([\"session\", \"aid\"]).unique()\n    test_sessions = test_df[\"session\"].unique().sort().to_list()\n    \n    n_sessions = len(test_sessions)\n    n_chunks = (n_sessions + chunk_size - 1) // chunk_size\n    print(f\"Processing inference in {n_chunks} chunks of {chunk_size:,} sessions...\")\n    \n    all_top_20 = {\n        \"clicks\": [],\n        \"carts\": [],\n        \"orders\": []\n    }\n    \n    for chunk_idx in range(n_chunks):\n        start = chunk_idx * chunk_size\n        end = min(start + chunk_size, n_sessions)\n        chunk_sess = test_sessions[start:end]\n        \n        t0 = time.time()\n        \n        # Filter history cho chunk này\n        chunk_hist = test_history.filter(pl.col(\"session\").is_in(chunk_sess))\n        chunk_hist_src = test_df.filter(pl.col(\"session\").is_in(chunk_sess))\n        \n        # 2. Candidates\n        chunk_candidates = generate_candidates(\n            chunk_hist,\n            covisit_clicks,\n            covisit_cart_order,\n            covisit_buy2buy,\n            popular_items,\n            chunk_size=chunk_size,\n        )\n        \n        # 3. Features\n        chunk_featured = add_all_features(\n            chunk_candidates, \n            chunk_hist_src,\n            global_item_feats_all,\n            global_item_feats_7d,\n        )\n        \n        # 4. Predict\n        for pred_type in [\"clicks\", \"carts\", \"orders\"]:\n            feats = feature_names[pred_type]\n            for f in feats:\n                if f not in chunk_featured.columns:\n                    chunk_featured = chunk_featured.with_columns(pl.lit(0.0).alias(f))\n            \n            X = chunk_featured.select(feats).to_numpy()\n            scores = models[pred_type].predict(X)\n            \n            preds_type_df = chunk_featured.select([\"session\", \"candidate_aid\"]).with_columns(\n                pl.Series(\"score\", scores)\n            )\n            \n            top_20 = (\n                preds_type_df\n                .sort(\"score\", descending=True)\n                .group_by(\"session\", maintain_order=False)\n                .head(20)\n                .group_by(\"session\")\n                .agg(pl.col(\"candidate_aid\").alias(\"aids\"))\n            )\n            all_top_20[pred_type].append(top_20)\n            \n        elapsed = time.time() - t0\n        print(f\"  Chunk {chunk_idx+1}/{n_chunks}: {len(chunk_sess):,} sessions completed in {elapsed:.1f}s\")\n        \n        del chunk_hist, chunk_hist_src, chunk_candidates, chunk_featured\n        gc.collect()\n        \n    # 5. Format submission.csv (Optimized with Polars)\n    print(\"\\nFormatting submission (Optimized with Polars)...\")\n    t_fmt = time.time()\n    \n    submission_parts = []\n    for pred_type in [\"clicks\", \"carts\", \"orders\"]:\n        top_20_all = pl.concat(all_top_20[pred_type])\n        \n        # Biến đổi nhanh trên Polars: aids (list[int]) -> chuỗi string ngăn cách bằng dấu cách\n        # Tạo cột session_type bằng cách cộng chuỗi\n        top_20_all = top_20_all.with_columns([\n            pl.col(\"aids\").list.eval(pl.element().cast(pl.String)).list.join(\" \").alias(\"labels\"),\n            (pl.col(\"session\").cast(pl.String) + pl.lit(f\"_{pred_type}\")).alias(\"session_type\")\n        ]).select([\"session_type\", \"labels\"])\n        \n        submission_parts.append(top_20_all)\n        \n        del top_20_all\n        gc.collect()\n        \n    submission = pl.concat(submission_parts)\n    submission.write_csv(output_path)\n    \n    print(f\"\\n Submission saved: {output_path} | Formatting time: {time.time() - t_fmt:.1f}s\")\n    print(f\"  Rows: {submission.height:,}\")\n    print(f\"  Preview:\")\n    print(submission.head(6))\n    \n    return submission\n\n\n\n# --- Tạo submission ---\ntest_raw = load_raw_data(TEST_PARQUET)\nsubmission = create_submission(test_raw, train_full_history, models, feature_names, global_item_feats_all, global_item_feats_7d)\n\n\n# ══════════════════════════════════════════════════════════════════════════════\n# CELL 11: Summary\n# ══════════════════════════════════════════════════════════════════════════════\n\nprint(\"\\n\" + \"=\" * 60)\nprint(\"PIPELINE SUMMARY\")\nprint(\"=\" * 60)\nprint(f\"Data:\")\nprint(f\"  Train events : {train_events_count:,}\")\nprint(f\"  Val events   : {val_df.height:,}\")\nprint(f\"\\nCo-visitation Matrices:\")\nprint(f\"  Clicks     : {covisit_clicks.height:,} pairs\")\nprint(f\"  Cart-Order : {covisit_cart_order.height:,} pairs\")\nprint(f\"  Buy2Buy    : {covisit_buy2buy.height:,} pairs\")\nprint(f\"\\nModels trained: {list(models.keys())}\")\nprint(f\"\\nValidation Recall@20:\")\nfor t, r in val_recalls.items():\n    w = RECALL_WEIGHTS[t]\n    print(f\"  {t:8s}: {r:.5f} (weight={w:.2f}, contribution={w*r:.5f})\")\nweighted = sum(RECALL_WEIGHTS[t] * val_recalls[t] for t in val_recalls)\nprint(f\"  {'TOTAL':8s}: {weighted:.5f}\")\nprint(\"=\" * 60)\n\n\n# ══════════════════════════════════════════════════════════════════════════════\n# CELL 12: Save Artifacts & Demo Standalone Batch Inference\n# ══════════════════════════════════════════════════════════════════════════════\n\n# 1. Lưu các artifacts xuống thư mục /kaggle/working/artifacts\nartifacts_dir = WORKING_DIR / \"artifacts\"\nsave_pipeline_artifacts(\n    covisit_clicks=covisit_clicks,\n    covisit_cart_order=covisit_cart_order,\n    covisit_buy2buy=covisit_buy2buy,\n    popular_items=popular_items,\n    global_item_feats_all=global_item_feats_all,\n    global_item_feats_7d=global_item_feats_7d,\n    models=models,\n    feature_names=feature_names,\n    artifacts_dir=artifacts_dir\n)\n\n# 2. Demo chạy load lại và inference thử\nprint(\"\\n\" + \"=\" * 60)\nprint(\"DEMO: LOAD ARTIFACTS AND RUN STANDALONE BATCH INFERENCE\")\nprint(\"=\" * 60)\n\n# Load lại từ đĩa để chứng minh sự độc lập\n(\n    loaded_clicks,\n    loaded_cart_order,\n    loaded_buy2buy,\n    loaded_popular,\n    loaded_feats_all,\n    loaded_feats_7d,\n    loaded_models,\n    loaded_features\n) = load_pipeline_artifacts(artifacts_dir=artifacts_dir)\n\n# Lấy thử 3 session từ test set làm ví dụ\nsample_sessions = test_raw[\"session\"].unique().head(3).to_list()\nsample_df = test_raw.filter(pl.col(\"session\").is_in(sample_sessions))\n\nprint(f\"\\nRunning standalone inference for sample sessions: {sample_sessions}\")\npredictions = recommend_for_batch(\n    session_events=sample_df,\n    covisit_clicks=loaded_clicks,\n    covisit_cart_order=loaded_cart_order,\n    covisit_buy2buy=loaded_buy2buy,\n    popular_items=loaded_popular,\n    global_item_feats_all=loaded_feats_all,\n    global_item_feats_7d=loaded_feats_7d,\n    models=loaded_models,\n    feature_names=loaded_features\n)\n\n# In kết quả gợi ý\nfor sess, preds in predictions.items():\n    print(f\"\\nSuggestions for Session {sess}:\")\n    print(f\"  - Clicks (Top 5): {preds['clicks'][:5]} (Total suggestions: {len(preds['clicks'])})\")\n    print(f\"  - Carts  (Top 5): {preds['carts'][:5]} (Total suggestions: {len(preds['carts'])})\")\n    print(f\"  - Orders (Top 5): {preds['orders'][:5]} (Total suggestions: {len(preds['orders'])})\")\nprint(\"=\" * 60)\n\n","metadata":{"_uuid":"ab872740-a72e-4399-a218-5a9c7968f218","_cell_guid":"46f52713-56aa-4ced-be71-fe438bc1b405","trusted":true,"collapsed":false,"jupyter":{"outputs_hidden":false},"execution":{"iopub.status.busy":"2026-05-28T03:24:00.713455Z","iopub.execute_input":"2026-05-28T03:24:00.713621Z","iopub.status.idle":"2026-05-28T03:56:44.608765Z","shell.execute_reply.started":"2026-05-28T03:24:00.713608Z","shell.execute_reply":"2026-05-28T03:56:44.608405Z"}},"outputs":[],"execution_count":null}]}