{"metadata":{"kernelspec":{"language":"python","display_name":"Python 3","name":"python3"},"language_info":{"name":"python","version":"3.12.12","mimetype":"text/x-python","codemirror_mode":{"name":"ipython","version":3},"pygments_lexer":"ipython3","nbconvert_exporter":"python","file_extension":".py"},"kaggle":{"accelerator":"nvidiaTeslaT4","dataSources":[{"sourceType":"competition","sourceId":31254,"databundleVersionId":3103714}],"dockerImageVersionId":31328,"isInternetEnabled":true,"language":"python","sourceType":"notebook","isGpuEnabled":true}},"nbformat_minor":4,"nbformat":4,"cells":[{"cell_type":"markdown","source":"# **CHƯƠNG 1: Tiền sử lý dữ liệu**","metadata":{}},{"cell_type":"markdown","source":"## **Giai Đoạn 1.1: EDA + Data Integration**","metadata":{}},{"cell_type":"code","source":"# ============================================================\n# CELL 0: Kiểm tra GPU & cài thư viện cần thiết\n# ============================================================\nimport subprocess, sys\n\n# Kiểm tra GPU có sẵn không\nresult = subprocess.run(['nvidia-smi'], capture_output=True, text=True)\nprint(result.stdout[:500])\n\n# Thử import RAPIDS cuDF (có sẵn trên Kaggle GPU notebook)\ntry:\n    import cudf\n    import cuml\n    USE_GPU = True\n    print(\"RAPIDS cuDF available — chạy trên GPU\")\nexcept ImportError:\n    USE_GPU = False\n    print(\"cuDF không có — fallback sang Polars (CPU)\")\n\n# Cài Polars nếu chưa có\ntry:\n    import polars as pl\n    print(\"Polars available\")\nexcept ImportError:\n    subprocess.run([sys.executable, '-m', 'pip', 'install', 'polars', '-q'])\n    import polars as pl\n\nprint(f\"\\nStrategy: {'cuDF (GPU)' if USE_GPU else 'Polars (CPU)'}\")","metadata":{"trusted":true,"execution":{"iopub.status.busy":"2026-04-06T16:24:49.783846Z","iopub.execute_input":"2026-04-06T16:24:49.784510Z","iopub.status.idle":"2026-04-06T16:24:58.101714Z","shell.execute_reply.started":"2026-04-06T16:24:49.784480Z","shell.execute_reply":"2026-04-06T16:24:58.100928Z"}},"outputs":[],"execution_count":null},{"cell_type":"code","source":"import rmm\nimport cudf\n\nrmm.reinitialize(\n    managed_memory=True, # Allows spilling to system memory\n    initial_pool_size=None \n)\n\nprint(\"Done!\")","metadata":{"trusted":true,"execution":{"iopub.status.busy":"2026-04-06T16:24:58.103105Z","iopub.execute_input":"2026-04-06T16:24:58.103616Z","iopub.status.idle":"2026-04-06T16:24:58.562934Z","shell.execute_reply.started":"2026-04-06T16:24:58.103585Z","shell.execute_reply":"2026-04-06T16:24:58.562315Z"}},"outputs":[],"execution_count":null},{"cell_type":"code","source":"# ============================================================\n# CELL 1: Cấu hình toàn cục\n# ============================================================\nimport os\nimport warnings\nimport gc\nwarnings.filterwarnings('ignore')\n\n# Đường dẫn trên Kaggle\nDATA_DIR = '/kaggle/input/competitions/h-and-m-personalized-fashion-recommendations'\nOUTPUT_DIR = '/kaggle/working'\n\nPATHS = {\n    'transactions': os.path.join(DATA_DIR, 'transactions_train.csv'),\n    'articles':     os.path.join(DATA_DIR, 'articles.csv'),\n    'customers':    os.path.join(DATA_DIR, 'customers.csv'),\n}\n\n# Kiểm tra file tồn tại\nfor name, path in PATHS.items():\n    size_mb = os.path.getsize(path) / 1e6\n    print(f\"  {'Yes'if os.path.exists(path) else 'No'} {name:15s}: {size_mb:8.1f} MB\")","metadata":{"trusted":true,"execution":{"iopub.status.busy":"2026-04-06T16:24:58.563750Z","iopub.execute_input":"2026-04-06T16:24:58.564036Z","iopub.status.idle":"2026-04-06T16:24:58.571202Z","shell.execute_reply.started":"2026-04-06T16:24:58.564011Z","shell.execute_reply":"2026-04-06T16:24:58.570419Z"}},"outputs":[],"execution_count":null},{"cell_type":"code","source":"# ============================================================\n# CELL 2: BƯỚC 1A — EDA transactions\n# BIỆN LUẬN: Đây là file lớn nhất (~3.7M dòng, ~500MB).\n# Dùng cuDF để load thẳng lên VRAM GPU, tránh OOM trên RAM CPU.\n# Trong bán lẻ thời trang, transaction là trung tâm: mọi\n# recommendation đều xuất phát từ hành vi mua thực tế.\n# ============================================================\n\ndef load_transactions(use_gpu=USE_GPU):\n    \"\"\"Load transactions với strategy tối ưu theo hardware.\"\"\"\n    \n    print(\"=\" * 60)\n    print(\"BƯỚC 1A: EDA TRANSACTIONS\")\n    print(\"=\" * 60)\n    \n    if use_gpu:\n        # cuDF load thẳng lên GPU VRAM\n        df = cudf.read_csv(\n            PATHS['transactions'],\n            dtype={\n                'customer_id': 'str',\n                'article_id':  'int32',   # 9 chữ số → int32 đủ\n                'price':       'float32', # float32 tiết kiệm 50% so với float64\n                'sales_channel_id': 'int8',\n            }\n        )\n        print(f\"  Backend: cuDF (GPU)\")\n    else:\n        # Polars lazy evaluation — chỉ đọc metadata trước\n        df = pl.read_csv(\n            PATHS['transactions'],\n            dtypes={\n                'customer_id': pl.Utf8,\n                'article_id':  pl.Int32,\n                'price':       pl.Float32,\n                'sales_channel_id': pl.Int8,\n            }\n        )\n        print(f\"  Backend: Polars (CPU)\")\n    \n    # --- Thông tin cơ bản ---\n    print(f\"\\n  Shape        : {df.shape}\")\n    print(f\"  Columns      : {list(df.columns)}\")\n    \n    # dtypes\n    print(f\"\\n  Dtypes:\")\n    if use_gpu:\n        for col, dtype in df.dtypes.items():\n            print(f\"    {col:30s}: {dtype}\")\n    else:\n        for col, dtype in zip(df.columns, df.dtypes):\n            print(f\"    {col:30s}: {dtype}\")\n    \n    # Null counts\n    print(f\"\\n  Null counts:\")\n    null_counts = df.isnull().sum() if use_gpu else df.null_count()\n    print(null_counts)\n    \n    # Thống kê mô tả\n    print(f\"\\n  Describe (price & sales_channel_id):\")\n    if use_gpu:\n        print(df[['price', 'sales_channel_id']].describe())\n    else:\n        print(df[['price', 'sales_channel_id']].describe())\n    \n    # Date range\n    if use_gpu:\n        print(f\"\\n  Date range: {df['t_dat'].min()} → {df['t_dat'].max()}\")\n    else:\n        print(f\"\\n  Date range: {df['t_dat'].min()} → {df['t_dat'].max()}\")\n    \n    # Unique counts — quan trọng để hiểu scale\n    print(f\"\\n  Unique customers : {df['customer_id'].nunique():,}\")\n    print(f\"  Unique articles  : {df['article_id'].nunique():,}\")\n    print(f\"  Total rows       : {len(df):,}\")\n    \n    # Phân phối sales_channel (1=online, 2=store)\n    print(f\"\\n  Sales channel distribution:\")\n    if use_gpu:\n        print(df['sales_channel_id'].value_counts())\n    else:\n        print(df['sales_channel_id'].value_counts())\n    \n    return df\n\ndf_trans = load_transactions()","metadata":{"trusted":true,"execution":{"iopub.status.busy":"2026-04-06T16:24:58.572671Z","iopub.execute_input":"2026-04-06T16:24:58.572874Z","iopub.status.idle":"2026-04-06T16:25:13.233259Z","shell.execute_reply.started":"2026-04-06T16:24:58.572856Z","shell.execute_reply":"2026-04-06T16:25:13.232411Z"}},"outputs":[],"execution_count":null},{"cell_type":"code","source":"# ============================================================\n# CELL 3: BƯỚC 1B — EDA articles\n# BIỆN LUẬN: articles.csv chứa metadata sản phẩm (~105K rows).\n# Cột product_type_name, colour_group_name sẽ là feature quan\n# trọng để nhận dạng xu hướng thời trang theo mùa.\n# ============================================================\n\ndef load_articles(use_gpu=USE_GPU):\n    print(\"=\" * 60)\n    print(\"BƯỚC 1B: EDA ARTICLES\")\n    print(\"=\" * 60)\n    \n    if use_gpu:\n        df = cudf.read_csv(\n            PATHS['articles'],\n            dtype={'article_id': 'int32'}\n        )\n    else:\n        df = pl.read_csv(PATHS['articles'])\n    \n    print(f\"  Shape   : {df.shape}\")\n    print(f\"  Columns : {list(df.columns)}\")\n    \n    # Null analysis — quan trọng cho feature engineering sau\n    print(f\"\\n  Null counts per column:\")\n    null_counts = df.isnull().sum() if use_gpu else df.null_count()\n    print(null_counts)\n    \n    # Phân phối product type — dùng để phân tích xu hướng\n    print(f\"\\n  Top 10 product_type_name:\")\n    if use_gpu:\n        print(df['product_type_name'].value_counts().head(10))\n    else:\n        print(\n            df.group_by('product_type_name')\n            .agg(pl.count().alias('count'))\n            .sort('count', descending=True)\n            .head(10)\n        )\n    \n    # Phân phối màu sắc\n    print(f\"\\n  Top 10 colour_group_name:\")\n    if use_gpu:\n        print(df['colour_group_name'].value_counts().head(10))\n    else:\n        print(\n            df.group_by('colour_group_name')\n            .agg(pl.count().alias('count'))\n            .sort('count', descending=True)\n            .head(10)\n        )\n    \n    # Department distribution\n    print(f\"\\n  Department distribution:\")\n    if use_gpu:\n        print(df['department_name'].value_counts().head(15))\n    else:\n        print(\n            df.group_by('department_name')\n            .agg(pl.count().alias('count'))\n            .sort('count', descending=True)\n            .head(15)\n        )\n    \n    # Sample rows\n    print(f\"\\n  Sample rows (5):\")\n    if use_gpu:\n        print(df.head(5).to_pandas().to_string())\n    else:\n        print(df.head(5))\n    \n    return df\n\ndf_art = load_articles()","metadata":{"trusted":true,"execution":{"iopub.status.busy":"2026-04-06T16:25:13.234327Z","iopub.execute_input":"2026-04-06T16:25:13.234748Z","iopub.status.idle":"2026-04-06T16:25:13.505514Z","shell.execute_reply.started":"2026-04-06T16:25:13.234714Z","shell.execute_reply":"2026-04-06T16:25:13.504854Z"}},"outputs":[],"execution_count":null},{"cell_type":"code","source":"# ============================================================\n# CELL 4: BƯỚC 1C — EDA customers\n# BIỆN LUẬN: customers.csv chứa ~1.37M user profiles.\n# Hai cột quan trọng: Age (phân khúc khách hàng) và\n# FN/Active/club_member_status (hành vi loyalty).\n# Null trong Age sẽ cần xử lý kỹ ở Bước 3.\n# ============================================================\n\ndef load_customers(use_gpu=USE_GPU):\n    print(\"=\" * 60)\n    print(\"BƯỚC 1C: EDA CUSTOMERS\")\n    print(\"=\" * 60)\n    \n    if use_gpu:\n        df = cudf.read_csv(PATHS['customers'])\n    else:\n        df = pl.read_csv(PATHS['customers'])\n    \n    print(f\"  Shape   : {df.shape}\")\n    print(f\"  Columns : {list(df.columns)}\")\n    \n    # Null analysis — cực kỳ quan trọng cho chương 2 bước 3\n    print(f\"\\n  Null analysis (%):\")\n    if use_gpu:\n        null_pct = (df.isnull().sum() / len(df) * 100).round(2)\n        print(null_pct)\n    else:\n        null_pct = df.null_count() / len(df) * 100\n        print(null_pct)\n    \n    # Age distribution\n    print(f\"\\n  Age stats:\")\n    if use_gpu:\n        print(df['age'].describe())\n    else:\n        print(df.select('age').describe())\n    \n    # Club member status\n    print(f\"\\n  club_member_status:\")\n    if use_gpu:\n        print(df['club_member_status'].value_counts())\n    else:\n        print(\n            df.group_by('club_member_status')\n            .agg(pl.count().alias('count'))\n            .sort('count', descending=True)\n        )\n    \n    # Fashion News Frequency\n    print(f\"\\n  fashion_news_frequency:\")\n    if use_gpu:\n        print(df['fashion_news_frequency'].value_counts())\n    else:\n        print(\n            df.group_by('fashion_news_frequency')\n            .agg(pl.count().alias('count'))\n            .sort('count', descending=True)\n        )\n    \n    # Active status\n    print(f\"\\n  Active column sample:\")\n    if use_gpu:\n        print(df['Active'].value_counts(dropna=False))\n    else:\n        print(df['Active'].value_counts())\n    \n    return df\n\ndf_cust = load_customers()","metadata":{"trusted":true,"execution":{"iopub.status.busy":"2026-04-06T16:25:13.506419Z","iopub.execute_input":"2026-04-06T16:25:13.506677Z","iopub.status.idle":"2026-04-06T16:25:14.571957Z","shell.execute_reply.started":"2026-04-06T16:25:13.506654Z","shell.execute_reply":"2026-04-06T16:25:14.571302Z"}},"outputs":[],"execution_count":null},{"cell_type":"code","source":"# ============================================================\n# CELL 5: BƯỚC 2A — Merge transactions + articles\n# BIỆN LUẬN: LEFT JOIN từ transactions → articles vì\n# chúng ta muốn GIỮ mọi transaction. Nếu có article_id\n# không match (data quality issue), ta phát hiện được thay\n# vì mất dòng. Trong bán lẻ, mất transaction = mất signal.\n# ============================================================\n\nimport time\n\ndef merge_transactions_articles(df_trans, df_art, use_gpu=USE_GPU):\n    print(\"=\" * 60)\n    print(\"BƯỚC 2A: MERGE transactions + articles\")\n    print(\"=\" * 60)\n    \n    t0 = time.time()\n    \n    # Chọn cột cần thiết từ articles để tránh bùng nổ RAM\n    # BIỆN LUẬN: articles có 25 cột nhưng ta chỉ cần\n    # metadata thời trang phục vụ feature engineering.\n    art_cols = [\n        'article_id',\n        'product_type_name',\n        'product_group_name',\n        'colour_group_name',\n        'department_name',\n        'section_name',\n        'garment_group_name',\n        'perceived_colour_master_name',\n    ]\n    \n    if use_gpu:\n        df_art_slim = df_art[art_cols]\n        df_merged = df_trans.merge(\n            df_art_slim,\n            on='article_id',\n            how='left',\n            sort=False  # không sort để tối ưu tốc độ GPU\n        )\n    else:\n        df_art_slim = df_art.select(art_cols)\n        df_merged = df_trans.join(\n            df_art_slim,\n            on='article_id',\n            how='left'\n        )\n    \n    elapsed = time.time() - t0\n    print(f\"  Merge xong trong {elapsed:.2f}s\")\n    print(f\"  Shape sau merge: {df_merged.shape}\")\n    \n    # Kiểm tra unmatched articles (data quality)\n    if use_gpu:\n        n_null = df_merged['product_type_name'].isnull().sum()\n    else:\n        n_null = df_merged['product_type_name'].null_count()\n    \n    print(f\"\\n Warning: Unmatched article_id: {n_null:,} rows ({n_null/len(df_merged)*100:.3f}%)\")\n    \n    if n_null > 0:\n        print(\"  → Sẽ xử lý ở Bước 3 (Data Cleaning)\")\n    \n    # Memory usage\n    if use_gpu:\n        mem_mb = df_merged.memory_usage(deep=True).sum() / 1e6\n        print(f\"\\n  GPU VRAM used: ~{mem_mb:.1f} MB\")\n    \n    print(f\"\\n  Sample (3 rows):\")\n    if use_gpu:\n        print(df_merged.head(3).to_pandas().to_string())\n    else:\n        print(df_merged.head(3))\n    \n    return df_merged\n\ndf_merged = merge_transactions_articles(df_trans, df_art)\n\n# Giải phóng memory\ndel df_art\ngc.collect()\nprint(\"\\n Cleaned df_art released from memory\")","metadata":{"trusted":true,"execution":{"iopub.status.busy":"2026-04-06T16:25:14.572934Z","iopub.execute_input":"2026-04-06T16:25:14.573196Z","iopub.status.idle":"2026-04-06T16:25:17.673183Z","shell.execute_reply.started":"2026-04-06T16:25:14.573171Z","shell.execute_reply":"2026-04-06T16:25:17.672608Z"}},"outputs":[],"execution_count":null},{"cell_type":"code","source":"import time\nimport gc\n\n# ============================================================\n# CELL 6: BƯỚC 2B — Merge với customers → df_master\n# ============================================================\n\ndef merge_with_customers(df_merged, df_cust, use_gpu=USE_GPU):\n    print(\"=\" * 60)\n    print(\"BƯỚC 2B: MERGE với customers → df_master\")\n    print(\"=\" * 60)\n    \n    t0 = time.time()\n    \n    # Chọn cột customers cần thiết\n    cust_cols = [\n        'customer_id',\n        'age',\n        'Active',\n        'club_member_status',\n        'fashion_news_frequency',\n        'FN',\n    ]\n    \n    if use_gpu:\n        # 1. Tạo bản slim và XÓA NGAY bản gốc lớn để trống VRAM\n        df_cust_slim = df_cust[cust_cols].copy()\n        del df_cust\n        gc.collect() \n        \n        # 2. Thực hiện Merge\n        df_master = df_merged.merge(\n            df_cust_slim,\n            on='customer_id',\n            how='left',\n            sort=False\n        )\n        \n        # 3. Dọn dẹp intermediate data ngay sau khi merge\n        del df_merged, df_cust_slim\n        gc.collect()\n        \n        elapsed = time.time() - t0\n        print(f\"  Merge xong trong {elapsed:.2f}s\")\n        print(f\"  Shape df_master: {df_master.shape}\")\n        \n        # 4. Null analysis cho GPU\n        print(f\"\\n Null summary df_master:\")\n        null_summary = df_master.isnull().sum()\n        null_pct = (null_summary / len(df_master) * 100).round(2)\n        for col in df_master.columns:\n            n = null_summary[col]\n            if n > 0:\n                print(f\"    {col:35s}: {n:>8,} ({null_pct[col]:.1f}%)\")\n                \n        # 5. Memory snapshot\n        mem_mb = df_master.memory_usage(deep=True).sum() / 1e6\n        print(f\"\\n  GPU VRAM used by df_master: ~{mem_mb:.1f} MB\")\n\n    else: # Chế độ CPU (Polars/Pandas)\n        df_cust_slim = df_cust.select(cust_cols)\n        df_master = df_merged.join(\n            df_cust_slim,\n            on='customer_id',\n            how='left'\n        )\n        \n        # Dọn dẹp\n        del df_cust, df_merged, df_cust_slim\n        gc.collect()\n        \n        elapsed = time.time() - t0\n        print(f\"  Merge xong trong {elapsed:.2f}s\")\n        print(f\"  Shape df_master: {df_master.shape}\")\n        \n        # Null analysis cho CPU\n        print(f\"\\n Null summary df_master:\")\n        null_counts = df_master.null_count()\n        for col in df_master.columns:\n            n = null_counts[col][0]\n            if n > 0:\n                pct = n / len(df_master) * 100\n                print(f\"    {col:35s}: {n:>8,} ({pct:.1f}%)\")\n\n    print(\"\\n Cleaned Intermediate DataFrames released\")\n    \n    return df_master\n\n# ============================================================\n# RUN KHỐI LỆNH\n# ============================================================\n\n# Ép kiểu Category TRƯỚC KHI GỌI HÀM để giảm thiểu tối đa bộ nhớ\ndf_merged['customer_id'] = df_merged['customer_id'].astype('category')\ndf_cust['customer_id'] = df_cust['customer_id'].astype('category')\n\n# Chỉ dùng hàm Merge (đã bỏ phần .map() để tránh trùng lặp)\ndf_master = merge_with_customers(df_merged, df_cust)","metadata":{"trusted":true,"execution":{"iopub.status.busy":"2026-04-06T16:25:17.674081Z","iopub.execute_input":"2026-04-06T16:25:17.674318Z","iopub.status.idle":"2026-04-06T16:25:21.463101Z","shell.execute_reply.started":"2026-04-06T16:25:17.674297Z","shell.execute_reply":"2026-04-06T16:25:21.462492Z"}},"outputs":[],"execution_count":null},{"cell_type":"code","source":"# ============================================================\n# CELL 7: EDA SUMMARY REPORT + SAVE SNAPSHOT\n# ============================================================\n\ndef eda_summary_report(df_master, use_gpu=USE_GPU):\n    print(\"=\" * 60)\n    print(\"EDA SUMMARY REPORT — df_master\")\n    print(\"=\" * 60)\n    \n    print(f\"\\n  Total records   : {len(df_master):>12,}\")\n    print(f\"  Total columns   : {df_master.shape[1]:>12,}\")\n    \n    # Unique counts\n    print(f\"\\n  Unique customers: {df_master['customer_id'].nunique():>12,}\")\n    print(f\"  Unique articles : {df_master['article_id'].nunique():>12,}\")\n    \n    # Date range\n    if use_gpu:\n        print(f\"  Date range      : {df_master['t_dat'].min()} → {df_master['t_dat'].max()}\")\n    else:\n        print(f\"  Date range      : {df_master['t_dat'].min()} → {df_master['t_dat'].max()}\")\n    \n    # Avg price overall\n    if use_gpu:\n        avg_price = float(df_master['price'].mean())\n    else:\n        avg_price = df_master['price'].mean()\n    print(f\"  Avg price       : {avg_price:>12.4f}\")\n    \n    # Transactions per customer (distribution)\n    print(f\"\\n  Transactions per customer:\")\n    if use_gpu:\n        txn_per_cust = df_master.groupby('customer_id')['article_id'].count()\n        print(txn_per_cust.describe())\n    else:\n        txn_per_cust = (\n            df_master.group_by('customer_id')\n            .agg(pl.count('article_id').alias('n_purchases'))\n        )\n        print(txn_per_cust.describe())\n    \n    print(f\"\\n  NHẬN XÉT:\")\n    print(f\"  - df_master là nguồn duy nhất cho Chương 3 & 4\")\n    print(f\"  - Cần xử lý missing Age/Active trước khi build features (Bước 3)\")\n    print(f\"  - Cần downcast dtype để tiết kiệm VRAM trước khi FP-Growth (Bước 4)\")\n    \n    # Save parquet để dùng lại (nhanh hơn CSV rất nhiều)\n    save_path = f'{OUTPUT_DIR}/df_master_raw.parquet'\n    \n    if use_gpu:\n        # WORKAROUND: Chuyển tạm thời category thành string để lưu Parquet\n        cat_cols = df_master.select_dtypes(include=['category']).columns\n        for col in cat_cols:\n            df_master[col] = df_master[col].astype(str)\n            \n        df_master.to_parquet(save_path)\n        \n        # Ép kiểu lại thành category ngay lập tức để tiết kiệm VRAM cho các bước sau\n        for col in cat_cols:\n            df_master[col] = df_master[col].astype('category')\n    else:\n        df_master.write_parquet(save_path)\n    \n    print(f\"\\n  Saved to: {save_path}\")\n    \n    import os\n    size_mb = os.path.getsize(save_path) / 1e6\n    print(f\"  File size: {size_mb:.1f} MB\")\n\neda_summary_report(df_master)","metadata":{"trusted":true,"execution":{"iopub.status.busy":"2026-04-06T16:25:21.464018Z","iopub.execute_input":"2026-04-06T16:25:21.464221Z","iopub.status.idle":"2026-04-06T16:25:28.194946Z","shell.execute_reply.started":"2026-04-06T16:25:21.464202Z","shell.execute_reply":"2026-04-06T16:25:28.194322Z"}},"outputs":[],"execution_count":null},{"cell_type":"code","source":"# ============================================================\n# CELL 8: VISUALISATION — 4 biểu đồ EDA cơ bản\n# ============================================================\nimport matplotlib.pyplot as plt\nimport matplotlib.ticker as mticker\nimport numpy as np\n\nfig, axes = plt.subplots(2, 2, figsize=(16, 10))\nfig.suptitle('H&M Dataset — EDA Overview (Chương 2, Bước 1-2)', \n             fontsize=14, fontweight='bold', y=1.01)\n\n# Helper: convert to pandas/numpy for matplotlib\ndef to_numpy(series):\n    \"\"\"Chuyển cuDF/Polars series sang numpy.\"\"\"\n    try:\n        return series.to_pandas().values\n    except:\n        return series.to_numpy()\n\n# ── Plot 1: Số giao dịch theo tháng ──\nax = axes[0, 0]\nif USE_GPU:\n    df_plot = df_master.copy()\n    df_plot['month'] = cudf.to_datetime(df_plot['t_dat']).dt.month\n    monthly = df_plot.groupby('month').size().to_pandas()\nelse:\n    df_plot = df_master.with_columns(\n        pl.col('t_dat').str.to_date().dt.month().alias('month')\n    )\n    monthly = (\n        df_plot.group_by('month')\n        .agg(pl.count().alias('n'))\n        .sort('month')\n        .to_pandas()\n    )\n    monthly = monthly.set_index('month')['n']\n\nax.bar(monthly.index, monthly.values, color='#3266ad', alpha=0.8)\nax.set_title('Giao dịch theo tháng')\nax.set_xlabel('Tháng')\nax.set_ylabel('Số giao dịch')\nax.yaxis.set_major_formatter(mticker.FuncFormatter(lambda x, _: f'{x/1e6:.1f}M'))\n\n# ── Plot 2: Phân phối Age ──\nax = axes[0, 1]\nif USE_GPU:\n    ages = df_master['age'].dropna().to_pandas()\nelse:\n    ages = df_master['age'].drop_nulls().to_numpy()\nax.hist(ages, bins=50, color='#e67e22', alpha=0.8, edgecolor='white', linewidth=0.5)\nax.set_title('Phân phối Age khách hàng')\nax.set_xlabel('Tuổi')\nax.set_ylabel('Số lượng')\nax.axvline(np.median(ages), color='red', linestyle='--', linewidth=1.5, label=f'Median: {np.median(ages):.0f}')\nax.legend(fontsize=9)\n\n# ── Plot 3: Top 10 product_type_name ──\nax = axes[1, 0]\nif USE_GPU:\n    top_products = (\n        df_master['product_type_name']\n        .value_counts()\n        .head(10)\n        .to_pandas()\n    )\n    # SỬA Ở ĐÂY: Lấy tên từ index, lấy số lượng từ values\n    labels = top_products.index.tolist()\n    values = top_products.values.tolist()\nelse:\n    top_products = (\n        df_master.group_by('product_type_name')\n        .agg(pl.count().alias('count'))\n        .sort('count', descending=True)\n        .head(10)\n        .to_pandas()\n    )\n    labels = top_products['product_type_name'].tolist()\n    values = top_products['count'].tolist()\n\ncolors = plt.cm.Blues(np.linspace(0.4, 0.9, len(labels)))\nbars = ax.barh(labels[::-1], values[::-1], color=colors[::-1])\nax.set_title('Top 10 Product Types')\nax.set_xlabel('Số giao dịch')\nax.xaxis.set_major_formatter(mticker.FuncFormatter(lambda x, _: f'{x/1e3:.0f}K'))\n\n# ── Plot 4: Phân phối giá ──\nax = axes[1, 1]\nif USE_GPU:\n    prices = df_master['price'].dropna().to_pandas()\nelse:\n    prices = df_master['price'].drop_nulls().to_numpy()\n# Clip outlier để biểu đồ dễ đọc\nprices_clipped = np.clip(prices, 0, np.percentile(prices, 99))\nax.hist(prices_clipped, bins=60, color='#27ae60', alpha=0.8, edgecolor='white', linewidth=0.3)\nax.set_title('Phân phối giá sản phẩm (clip tại P99)')\nax.set_xlabel('Price')\nax.set_ylabel('Số lượng')\nax.axvline(np.median(prices_clipped), color='red', linestyle='--', linewidth=1.5,\n           label=f'Median: {np.median(prices_clipped):.4f}')\nax.legend(fontsize=9)\n\nplt.tight_layout()\nplt.savefig(f'{OUTPUT_DIR}/eda_overview.png', dpi=150, bbox_inches='tight')\nplt.show()\nprint(\"EDA plot saved!\")","metadata":{"trusted":true,"execution":{"iopub.status.busy":"2026-04-06T16:25:28.197090Z","iopub.execute_input":"2026-04-06T16:25:28.197503Z","iopub.status.idle":"2026-04-06T16:25:35.019899Z","shell.execute_reply.started":"2026-04-06T16:25:28.197480Z","shell.execute_reply":"2026-04-06T16:25:35.019116Z"}},"outputs":[],"execution_count":null},{"cell_type":"markdown","source":"## **Giai Đoạn 1.2: Data Cleaning + Downcasting**","metadata":{}},{"cell_type":"code","source":"# ============================================================\n# CELL 9: BƯỚC 3A — Xử lý missing Age (2 cách + biện luận)\n#\n# BIỆN LUẬN: Age là feature phân khúc khách hàng quan trọng\n# trong bán lẻ thời trang (teen vs. adult vs. senior mua rất\n# khác nhau). Tuy nhiên ~20% null là nhiều, nên KHÔNG drop —\n# ta so sánh 2 chiến lược và chọn cách phù hợp học thuật.\n# ============================================================\nimport numpy as np\n\ndef analyze_missing_age(df_master, use_gpu=USE_GPU):\n    print(\"=\" * 60)\n    print(\"BƯỚC 3A: MISSING AGE — Phân tích\")\n    print(\"=\" * 60)\n\n    if use_gpu:\n        n_null = int(df_master['age'].isnull().sum())\n        total  = len(df_master)\n        age_vals = df_master['age'].dropna().to_pandas()\n    else:\n        n_null = df_master['age'].null_count()\n        total  = len(df_master)\n        age_vals = df_master['age'].drop_nulls().to_numpy()\n\n    print(f\"  Null count : {n_null:,} / {total:,} ({n_null/total*100:.1f}%)\")\n    print(f\"  Median age : {np.median(age_vals):.1f}\")\n    print(f\"  Mean age   : {np.mean(age_vals):.1f}\")\n    print(f\"  Std age    : {np.std(age_vals):.1f}\")\n    print(f\"  Range      : {np.min(age_vals):.0f} – {np.max(age_vals):.0f}\")\n\nanalyze_missing_age(df_master)\n\n\ndef apply_age_strategy(df_master, strategy='median', use_gpu=USE_GPU):\n    \"\"\"\n    strategy='median' → Cách 1: điền median\n    strategy='flag'   → Cách 2: giữ null, tạo cột age_group_known\n    \n    BIỆN LUẬN CÁC CHIẾN LƯỢC:\n    \n    Cách 1 - Median Imputation:\n      Ưu: Đơn giản, không bias phân phối, median robust hơn mean\n          với outliers tuổi (VD: 16 vs 90).\n      Nhược: Xoá đi tín hiệu \"unknown age\" — trong thực tế,\n          khách không cung cấp tuổi thường là nhóm trẻ hơn.\n    \n    Cách 2 - Flag + Keep:\n      Ưu: Bảo toàn thông tin \"null\" như một đặc trưng riêng.\n          Mô hình có thể học: null_age → hành vi mua khác.\n      Nhược: Cần handle NaN khi tính feature (phức tạp hơn).\n    \n    → QUYẾT ĐỊNH: Dùng Cách 2 cho ML pipeline (Chương 4),\n      Cách 1 cho Association Rules (FP-Growth không xử lý null).\n    \"\"\"\n    print(f\"\\n  Áp dụng strategy: '{strategy}'\")\n\n    if strategy == 'median':\n        if use_gpu:\n            median_val = float(df_master['age'].median())\n            df_master['age'] = df_master['age'].fillna(median_val)\n            df_master['age_imputed'] = (df_master['age'] == median_val).astype('int8')\n        else:\n            median_val = df_master['age'].median()\n            df_master = df_master.with_columns([\n                pl.col('age').fill_null(median_val).alias('age'),\n                pl.when(pl.col('age').is_null())\n                  .then(pl.lit(1))\n                  .otherwise(pl.lit(0))\n                  .cast(pl.Int8)\n                  .alias('age_imputed')\n            ])\n        print(f\"  → Filled with median = {median_val:.1f}\")\n        print(f\"  → Thêm cột 'age_imputed' (1 = was null)\")\n\n    elif strategy == 'flag':\n        # Tạo age_group: bin tuổi thành nhóm + nhóm Unknown\n        if use_gpu:\n            import cudf\n            age = df_master['age']\n            conditions = [\n                age < 25,\n                (age >= 25) & (age < 35),\n                (age >= 35) & (age < 50),\n                age >= 50,\n            ]\n            choices = ['Teen/GenZ', 'Millennial', 'GenX', 'Boomer+']\n            # cuDF: dùng nested where\n            df_master['age_group'] = 'Unknown'\n            for cond, label in zip(conditions, choices):\n                df_master['age_group'] = df_master['age_group'].where(~cond, label)\n        else:\n            df_master = df_master.with_columns(\n                pl.when(pl.col('age').is_null()).then(pl.lit('Unknown'))\n                  .when(pl.col('age') < 25).then(pl.lit('Teen/GenZ'))\n                  .when(pl.col('age') < 35).then(pl.lit('Millennial'))\n                  .when(pl.col('age') < 50).then(pl.lit('GenX'))\n                  .otherwise(pl.lit('Boomer+'))\n                  .alias('age_group')\n            )\n        print(f\"  → Tạo cột 'age_group' (5 bins + Unknown)\")\n        if use_gpu:\n            print(df_master['age_group'].value_counts())\n        else:\n            print(\n                df_master.group_by('age_group')\n                .agg(pl.count().alias('count'))\n                .sort('count', descending=True)\n            )\n\n    return df_master\n\n# Áp dụng cả 2 chiến lược trên df_master\n# Thực tế: ta dùng 'flag' làm primary, 'median' bản copy cho FP-Growth\ndf_master = apply_age_strategy(df_master, strategy='flag')\n\n# Tạo bản sao age đã điền median để dùng sau\nif USE_GPU:\n    median_age = float(df_master['age'].median())\n    df_master['age_filled'] = df_master['age'].fillna(median_age).astype('float32')\nelse:\n    median_age = df_master['age'].median()\n    df_master = df_master.with_columns(\n        pl.col('age').fill_null(median_age).cast(pl.Float32).alias('age_filled')\n    )\n\nprint(f\"\\n  age_filled (median={median_age:.1f}) + age_group (flag) — cả 2 ready\")","metadata":{"trusted":true,"execution":{"iopub.status.busy":"2026-04-06T16:25:35.020780Z","iopub.execute_input":"2026-04-06T16:25:35.021061Z","iopub.status.idle":"2026-04-06T16:25:37.181929Z","shell.execute_reply.started":"2026-04-06T16:25:35.021028Z","shell.execute_reply":"2026-04-06T16:25:37.181277Z"}},"outputs":[],"execution_count":null},{"cell_type":"code","source":"# ============================================================\n# CELL 10: BƯỚC 3B — Missing Active, FN, club_member_status\n#\n# BIỆN LUẬN: Cột 'Active' và 'FN' là binary signal loyalty.\n# Khách hàng null ở đây thường là anonymous/guest purchase —\n# ta KHÔNG drop mà gán flag riêng (-1) để mô hình phân biệt\n# được 3 trạng thái: Active=1, Inactive=0, Unknown=-1.\n# Trong bán lẻ, \"không biết\" khác với \"không active\".\n# ============================================================\n\ndef clean_binary_flags(df_master, use_gpu=USE_GPU):\n    print(\"=\" * 60)\n    print(\"BƯỚC 3B: MISSING BINARY FLAGS\")\n    print(\"=\" * 60)\n\n    binary_cols = ['Active', 'FN']\n\n    for col in binary_cols:\n        if use_gpu:\n            n_null = int(df_master[col].isnull().sum())\n        else:\n            n_null = df_master[col].null_count()\n        pct = n_null / len(df_master) * 100\n        print(f\"\\n  {col}: {n_null:,} nulls ({pct:.1f}%)\")\n\n    # Strategy: fillna(-1) → 3-state encoding\n    # Active: 1=active, 0=inactive, -1=unknown\n    if use_gpu:\n        for col in binary_cols:\n            df_master[col] = df_master[col].fillna(-1).astype('int8')\n    else:\n        fill_exprs = [\n            pl.col(col).fill_null(-1).cast(pl.Int8).alias(col)\n            for col in binary_cols\n        ]\n        df_master = df_master.with_columns(fill_exprs)\n\n    print(f\"\\n  → Binary flags filled with -1 (Unknown state)\")\n    if use_gpu:\n        for col in binary_cols:\n            print(f\"  {col}: {df_master[col].value_counts()}\")\n    else:\n        for col in binary_cols:\n            print(f\"\\n  {col}:\")\n            print(df_master.group_by(col).agg(pl.count().alias('n')).sort('n', descending=True))\n\n    # club_member_status: categorical → fill 'NOT_MEMBER'\n    if use_gpu:\n        n_null_club = int(df_master['club_member_status'].isnull().sum())\n        df_master['club_member_status'] = df_master['club_member_status'].fillna('NOT_MEMBER')\n    else:\n        n_null_club = df_master['club_member_status'].null_count()\n        df_master = df_master.with_columns(\n            pl.col('club_member_status').fill_null('NOT_MEMBER')\n        )\n    print(f\"\\n  club_member_status: {n_null_club:,} nulls → filled 'NOT_MEMBER'\")\n\n    # fashion_news_frequency: fill 'NONE'\n    if use_gpu:\n        df_master['fashion_news_frequency'] = df_master['fashion_news_frequency'].fillna('NONE')\n    else:\n        df_master = df_master.with_columns(\n            pl.col('fashion_news_frequency').fill_null('NONE')\n        )\n\n    # Unmatched articles (từ merge bước 2) → fill 'UNKNOWN'\n    cat_article_cols = ['product_type_name', 'product_group_name',\n                        'colour_group_name', 'department_name',\n                        'section_name', 'garment_group_name',\n                        'perceived_colour_master_name']\n    if use_gpu:\n        for col in cat_article_cols:\n            df_master[col] = df_master[col].fillna('UNKNOWN')\n    else:\n        fill_cat = [pl.col(c).fill_null('UNKNOWN') for c in cat_article_cols]\n        df_master = df_master.with_columns(fill_cat)\n\n    print(f\"  Article category nulls → filled 'UNKNOWN'\")\n\n    # Final null check\n    print(f\"\\n  Null check sau cleaning:\")\n    if use_gpu:\n        remaining = df_master.isnull().sum()\n        remaining = remaining[remaining > 0]\n        print(remaining if len(remaining) > 0 else \" Không còn null (ngoài age raw)\")\n    else:\n        null_counts = df_master.null_count()\n        has_null = {col: null_counts[col][0] for col in df_master.columns\n                    if null_counts[col][0] > 0}\n        print(has_null if has_null else \" Không còn null (ngoài age raw)\")\n\n    return df_master\n\ndf_master = clean_binary_flags(df_master)","metadata":{"trusted":true,"execution":{"iopub.status.busy":"2026-04-06T16:25:37.182900Z","iopub.execute_input":"2026-04-06T16:25:37.183192Z","iopub.status.idle":"2026-04-06T16:25:38.379161Z","shell.execute_reply.started":"2026-04-06T16:25:37.183168Z","shell.execute_reply":"2026-04-06T16:25:38.378484Z"}},"outputs":[],"execution_count":null},{"cell_type":"code","source":"# ============================================================\n# CELL 10: BƯỚC 3B — Missing Active, FN, club_member_status\n#\n# BIỆN LUẬN: Cột 'Active' và 'FN' là binary signal loyalty.\n# Khách hàng null ở đây thường là anonymous/guest purchase —\n# ta KHÔNG drop mà gán flag riêng (-1) để mô hình phân biệt\n# được 3 trạng thái: Active=1, Inactive=0, Unknown=-1.\n# Trong bán lẻ, \"không biết\" khác với \"không active\".\n# ============================================================\n\ndef clean_binary_flags(df_master, use_gpu=USE_GPU):\n    print(\"=\" * 60)\n    print(\"BƯỚC 3B: MISSING BINARY FLAGS\")\n    print(\"=\" * 60)\n\n    binary_cols = ['Active', 'FN']\n\n    for col in binary_cols:\n        if use_gpu:\n            n_null = int(df_master[col].isnull().sum())\n        else:\n            n_null = df_master[col].null_count()\n        pct = n_null / len(df_master) * 100\n        print(f\"\\n  {col}: {n_null:,} nulls ({pct:.1f}%)\")\n\n    # Strategy: fillna(-1) → 3-state encoding\n    # Active: 1=active, 0=inactive, -1=unknown\n    if use_gpu:\n        for col in binary_cols:\n            df_master[col] = df_master[col].fillna(-1).astype('int8')\n    else:\n        fill_exprs = [\n            pl.col(col).fill_null(-1).cast(pl.Int8).alias(col)\n            for col in binary_cols\n        ]\n        df_master = df_master.with_columns(fill_exprs)\n\n    print(f\"\\n  → Binary flags filled with -1 (Unknown state)\")\n    if use_gpu:\n        for col in binary_cols:\n            print(f\"  {col}: {df_master[col].value_counts()}\")\n    else:\n        for col in binary_cols:\n            print(f\"\\n  {col}:\")\n            print(df_master.group_by(col).agg(pl.count().alias('n')).sort('n', descending=True))\n\n    # club_member_status: categorical → fill 'NOT_MEMBER'\n    if use_gpu:\n        n_null_club = int(df_master['club_member_status'].isnull().sum())\n        df_master['club_member_status'] = df_master['club_member_status'].fillna('NOT_MEMBER')\n    else:\n        n_null_club = df_master['club_member_status'].null_count()\n        df_master = df_master.with_columns(\n            pl.col('club_member_status').fill_null('NOT_MEMBER')\n        )\n    print(f\"\\n  club_member_status: {n_null_club:,} nulls → filled 'NOT_MEMBER'\")\n\n    # fashion_news_frequency: fill 'NONE'\n    if use_gpu:\n        df_master['fashion_news_frequency'] = df_master['fashion_news_frequency'].fillna('NONE')\n    else:\n        df_master = df_master.with_columns(\n            pl.col('fashion_news_frequency').fill_null('NONE')\n        )\n\n    # Unmatched articles (từ merge bước 2) → fill 'UNKNOWN'\n    cat_article_cols = ['product_type_name', 'product_group_name',\n                        'colour_group_name', 'department_name',\n                        'section_name', 'garment_group_name',\n                        'perceived_colour_master_name']\n    if use_gpu:\n        for col in cat_article_cols:\n            df_master[col] = df_master[col].fillna('UNKNOWN')\n    else:\n        fill_cat = [pl.col(c).fill_null('UNKNOWN') for c in cat_article_cols]\n        df_master = df_master.with_columns(fill_cat)\n\n    print(f\"  Article category nulls → filled 'UNKNOWN'\")\n\n    # Final null check\n    print(f\"\\n  Null check sau cleaning:\")\n    if use_gpu:\n        remaining = df_master.isnull().sum()\n        remaining = remaining[remaining > 0]\n        print(remaining if len(remaining) > 0 else \" Không còn null (ngoài age raw)\")\n    else:\n        null_counts = df_master.null_count()\n        has_null = {col: null_counts[col][0] for col in df_master.columns\n                    if null_counts[col][0] > 0}\n        print(has_null if has_null else \" Không còn null (ngoài age raw)\")\n\n    return df_master\n\ndf_master = clean_binary_flags(df_master)","metadata":{"trusted":true,"execution":{"iopub.status.busy":"2026-04-06T16:25:38.379990Z","iopub.execute_input":"2026-04-06T16:25:38.380259Z","iopub.status.idle":"2026-04-06T16:25:39.315161Z","shell.execute_reply.started":"2026-04-06T16:25:38.380222Z","shell.execute_reply":"2026-04-06T16:25:39.314521Z"}},"outputs":[],"execution_count":null},{"cell_type":"code","source":"# ============================================================\n# CELL 12: BƯỚC 4 — DOWNCASTING (Tối ưu RAM/VRAM)\n#\n# BIỆN LUẬN: Với ~31M dòng, mỗi cột float64 chiếm ~248MB.\n# Downcasting sang float32/int32/int8 có thể tiết kiệm 40-60%\n# VRAM — cực kỳ quan trọng để FP-Growth (Chương 4) chạy được\n# trên T4 (16GB VRAM mỗi GPU). Thời trang không cần precision\n# 64-bit cho giá trị price hay age.\n# ============================================================\n\ndef downcast_dtypes(df_master, use_gpu=USE_GPU):\n    print(\"=\" * 60)\n    print(\"BƯỚC 4: DOWNCASTING DTYPE\")\n    print(\"=\" * 60)\n\n    # Tính memory trước\n    if use_gpu:\n        mem_before = df_master.memory_usage(deep=True).sum() / 1e6\n    else:\n        mem_before = df_master.estimated_size('mb')\n    print(f\"  Memory TRƯỚC: {mem_before:.1f} MB\")\n\n    if use_gpu:\n        # --- INT columns ---\n        # article_id: max ~9 chữ số → int32 (max 2.1B) đủ\n        df_master['article_id'] = df_master['article_id'].astype('int32')\n\n        # sales_channel_id: chỉ 1 hoặc 2 → int8\n        df_master['sales_channel_id'] = df_master['sales_channel_id'].astype('int8')\n\n        # purchase_count: max vài nghìn → int16 (max 32K)\n        if 'purchase_count' in df_master.columns:\n            df_master['purchase_count'] = df_master['purchase_count'].clip(0, 32767).astype('int16')\n\n        # --- FLOAT columns ---\n        # price: float64 → float32 (6-7 chữ số thập phân đủ cho giá)\n        df_master['price'] = df_master['price'].astype('float32')\n        df_master['age_filled'] = df_master['age_filled'].astype('float32')\n\n        # age raw: float64 → float32\n        df_master['age'] = df_master['age'].astype('float32')\n\n        # --- INT8 flags (đã xử lý ở bước 3) ---\n        for col in ['Active', 'FN']:\n            if col in df_master.columns:\n                df_master[col] = df_master[col].astype('int8')\n\n        mem_after = df_master.memory_usage(deep=True).sum() / 1e6\n\n    else:\n        # Polars downcasting\n        cast_exprs = [\n            pl.col('article_id').cast(pl.Int32),\n            pl.col('sales_channel_id').cast(pl.Int8),\n            pl.col('price').cast(pl.Float32),\n            pl.col('age_filled').cast(pl.Float32),\n        ]\n        # Optional cols\n        if 'purchase_count' in df_master.columns:\n            cast_exprs.append(\n                pl.col('purchase_count').clip(0, 32767).cast(pl.Int16)\n            )\n        for col in ['Active', 'FN']:\n            if col in df_master.columns:\n                cast_exprs.append(pl.col(col).cast(pl.Int8))\n\n        df_master = df_master.with_columns(cast_exprs)\n        mem_after = df_master.estimated_size('mb')\n\n    savings = (mem_before - mem_after) / mem_before * 100\n    print(f\"  Memory SAU  : {mem_after:.1f} MB\")\n    print(f\"  Tiết kiệm   : {mem_before - mem_after:.1f} MB ({savings:.1f}%)\")\n\n    # In bảng dtypes sau downcast\n    print(f\"\\n  Dtype summary sau downcast:\")\n    if use_gpu:\n        for col, dtype in df_master.dtypes.items():\n            print(f\"    {col:40s}: {str(dtype):10s}\")\n    else:\n        for col, dtype in zip(df_master.columns, df_master.dtypes):\n            print(f\"    {col:40s}: {str(dtype):10s}\")\n\n    return df_master\n\ndf_master = downcast_dtypes(df_master)","metadata":{"trusted":true,"execution":{"iopub.status.busy":"2026-04-06T16:25:39.316177Z","iopub.execute_input":"2026-04-06T16:25:39.316473Z","iopub.status.idle":"2026-04-06T16:25:39.429064Z","shell.execute_reply.started":"2026-04-06T16:25:39.316440Z","shell.execute_reply":"2026-04-06T16:25:39.428479Z"}},"outputs":[],"execution_count":null},{"cell_type":"code","source":"# ============================================================\n# CELL 13: SAVE CHECKPOINT df_master_cleaned\n# ============================================================\nimport os\n\ndef save_cleaned_checkpoint(df_master, use_gpu=USE_GPU):\n    path = f'{OUTPUT_DIR}/df_master_cleaned.parquet'\n\n    if use_gpu:\n        # WORKAROUND: Tạm thời chuyển tất cả các cột category sang string để lưu\n        cat_cols = df_master.select_dtypes(include=['category']).columns\n        for col in cat_cols:\n            df_master[col] = df_master[col].astype(str)\n            \n        df_master.to_parquet(path)\n        \n        # Ép kiểu ngược lại ngay lập tức để giữ VRAM thấp cho các bước tiếp theo\n        for col in cat_cols:\n            df_master[col] = df_master[col].astype('category')\n    else:\n        df_master.write_parquet(path)\n\n    size_mb = os.path.getsize(path) / 1e6\n    print(f\"   Saved: {path}\")\n    print(f\"   Size: {size_mb:.1f} MB\")\n    print(f\"   Shape: {df_master.shape}\")\n    print(f\"   Columns: {list(df_master.columns)}\")\n\n    # Summary cleaning\n    print(f\"\\nCLEANING SUMMARY:\")\n    print(f\"  - age         : null → age_group (flag) + age_filled (median)\")\n    print(f\"  - Active/FN   : null → -1 (Unknown state)\")\n    print(f\"  - club/news   : null → 'NOT_MEMBER' / 'NONE'\")\n    print(f\"  - article cats: null → 'UNKNOWN'\")\n    print(f\"  - buyer_type  : IQR-based outlier tagging (không drop)\")\n    print(f\"  - dtype       : downcasted ~40-60% RAM savings\")\n\nsave_cleaned_checkpoint(df_master)","metadata":{"trusted":true,"execution":{"iopub.status.busy":"2026-04-06T16:25:39.429850Z","iopub.execute_input":"2026-04-06T16:25:39.430122Z","iopub.status.idle":"2026-04-06T16:25:46.732347Z","shell.execute_reply.started":"2026-04-06T16:25:39.430087Z","shell.execute_reply":"2026-04-06T16:25:46.731613Z"}},"outputs":[],"execution_count":null},{"cell_type":"markdown","source":"## **Giai Đoạn 1.3: Feature Engineering + Critical Thinking**","metadata":{}},{"cell_type":"code","source":"# ============================================================\n# CELL 14: BƯỚC 5A — TEMPORAL FEATURES\n#\n# BIỆN LUẬN: Thời trang có chu kỳ mùa rõ rệt (Spring/Summer\n# vs. Fall/Winter — SS/FW collection). Nếu không encode mùa,\n# mô hình sẽ không nhận ra pattern \"mua áo khoác vào Thu-Đông\".\n# Week-of-year quan trọng hơn month vì H&M drop collection\n# theo tuần (ví dụ: sale cuối tuần, launch thứ Năm).\n# ============================================================\nimport gc\n\ndef add_temporal_features(df_master, use_gpu=USE_GPU):\n    print(\"=\" * 60)\n    print(\"BƯỚC 5A: TEMPORAL FEATURES\")\n    print(\"=\" * 60)\n\n    if use_gpu:\n        import cudf\n\n        # Parse date (string → datetime)\n        df_master['t_dat'] = cudf.to_datetime(df_master['t_dat'])\n\n        # Cơ bản\n        df_master['year']         = df_master['t_dat'].dt.year.astype('int16')\n        df_master['month']        = df_master['t_dat'].dt.month.astype('int8')\n        df_master['week_of_year'] = df_master['t_dat'].dt.isocalendar().week.astype('int8')\n        df_master['day_of_week']  = df_master['t_dat'].dt.dayofweek.astype('int8')\n        df_master['is_weekend']   = (df_master['day_of_week'] >= 5).astype('int8')\n\n        # Mùa thời trang theo industry standard H&M:\n        # SS (Spring/Summer): tháng 2–7\n        # FW (Fall/Winter)  : tháng 8–1\n        month = df_master['month']\n        df_master['fashion_season'] = 'FW'  # default\n        df_master['fashion_season'] = df_master['fashion_season'].where(\n            ~((month >= 2) & (month <= 7)), 'SS'\n        )\n\n        # Fashion quarter (H&M drop cycle thực tế)\n        # Q1: Jan–Mar (FW clearance + SS preview)\n        # Q2: Apr–Jun (SS peak)\n        # Q3: Jul–Sep (SS sale + FW preview)\n        # Q4: Oct–Dec (FW peak + holiday)\n        df_master['fashion_quarter'] = 'Q4'\n        for q, (m_start, m_end) in enumerate([(1,3),(4,6),(7,9)], start=1):\n            cond = (month >= m_start) & (month <= m_end)\n            df_master['fashion_quarter'] = df_master['fashion_quarter'].where(\n                ~cond, f'Q{q}'\n            )\n\n        # Days since first purchase (recency anchor)\n        min_date = df_master['t_dat'].min()\n        df_master['days_since_start'] = (\n            (df_master['t_dat'] - min_date).dt.days.astype('int16')\n        )\n\n    else:\n        import polars as pl\n\n        df_master = df_master.with_columns(\n            pl.col('t_dat').str.to_date().alias('t_dat_parsed')\n        )\n\n        df_master = df_master.with_columns([\n            pl.col('t_dat_parsed').dt.year().cast(pl.Int16).alias('year'),\n            pl.col('t_dat_parsed').dt.month().cast(pl.Int8).alias('month'),\n            pl.col('t_dat_parsed').dt.week().cast(pl.Int8).alias('week_of_year'),\n            pl.col('t_dat_parsed').dt.weekday().cast(pl.Int8).alias('day_of_week'),\n            (pl.col('t_dat_parsed').dt.weekday() >= 5).cast(pl.Int8).alias('is_weekend'),\n        ])\n\n        df_master = df_master.with_columns(\n            pl.when((pl.col('month') >= 2) & (pl.col('month') <= 7))\n              .then(pl.lit('SS'))\n              .otherwise(pl.lit('FW'))\n              .alias('fashion_season')\n        )\n\n        df_master = df_master.with_columns(\n            pl.when((pl.col('month') >= 1) & (pl.col('month') <= 3)).then(pl.lit('Q1'))\n              .when((pl.col('month') >= 4) & (pl.col('month') <= 6)).then(pl.lit('Q2'))\n              .when((pl.col('month') >= 7) & (pl.col('month') <= 9)).then(pl.lit('Q3'))\n              .otherwise(pl.lit('Q4'))\n              .alias('fashion_quarter')\n        )\n\n        min_date = df_master['t_dat_parsed'].min()\n        df_master = df_master.with_columns(\n            (pl.col('t_dat_parsed') - min_date)\n            .dt.total_days().cast(pl.Int16)\n            .alias('days_since_start')\n        )\n\n    # Report\n    print(f\"  New columns: year, month, week_of_year, day_of_week,\")\n    print(f\"               is_weekend, fashion_season, fashion_quarter,\")\n    print(f\"               days_since_start\")\n\n    print(f\"\\n  Fashion season distribution:\")\n    if use_gpu:\n        print(df_master['fashion_season'].value_counts())\n    else:\n        print(\n            df_master.group_by('fashion_season')\n            .agg(pl.count().alias('n'))\n            .sort('n', descending=True)\n        )\n\n    print(f\"\\n  Fashion quarter distribution:\")\n    if use_gpu:\n        # SỬA Ở DÒNG NÀY: Dùng sort_index() thay vì sort_values()\n        print(df_master['fashion_quarter'].value_counts().sort_index())\n    else:\n        print(\n            df_master.group_by('fashion_quarter')\n            .agg(pl.count().alias('n'))\n            .sort('fashion_quarter')\n        )\n\n    print(f\"\\n  Weekend transactions:\")\n    if use_gpu:\n        n_wknd = int(df_master['is_weekend'].sum())\n        print(f\"  {n_wknd:,} ({n_wknd/len(df_master)*100:.1f}%) — Sat/Sun\")\n    else:\n        n_wknd = df_master['is_weekend'].sum()\n        print(f\"  {n_wknd:,} ({n_wknd/len(df_master)*100:.1f}%) — Sat/Sun\")\n\n    return df_master\n\ndf_master = add_temporal_features(df_master)\ngc.collect()\nprint(\"\\nTemporal features done\")","metadata":{"trusted":true,"execution":{"iopub.status.busy":"2026-04-06T16:25:46.733433Z","iopub.execute_input":"2026-04-06T16:25:46.733686Z","iopub.status.idle":"2026-04-06T16:25:49.044283Z","shell.execute_reply.started":"2026-04-06T16:25:46.733663Z","shell.execute_reply":"2026-04-06T16:25:49.043630Z"}},"outputs":[],"execution_count":null},{"cell_type":"code","source":"# ============================================================\n# CELL 15: BƯỚC 5B — PRICE FEATURES\n#\n# BIỆN LUẬN: Giá là tín hiệu phân khúc khách hàng mạnh nhất\n# trong thời trang. Khách budget buyer luôn mua dưới ngưỡng\n# giá trung bình, premium buyer mua trên — pattern này ổn định\n# và giúp recommendation tránh đề xuất sản phẩm \"lệch tầm\".\n# price_ratio (giá sp / avg giá KH) là feature tương đối, tốt\n# hơn giá tuyệt đối vì không bị lệch theo mùa sale.\n# ============================================================\n\ndef add_price_features(df_master, use_gpu=USE_GPU):\n    print(\"=\" * 60)\n    print(\"BƯỚC 5B: PRICE FEATURES\")\n    print(\"=\" * 60)\n\n    if use_gpu:\n        # Giá trung bình mỗi khách hàng\n        avg_price_cust = (\n            df_master.groupby('customer_id')['price']\n            .mean()\n            .reset_index()\n            .rename(columns={'price': 'avg_price_per_customer'})\n        )\n        df_master = df_master.merge(avg_price_cust, on='customer_id', how='left')\n        df_master['avg_price_per_customer'] = \\\n            df_master['avg_price_per_customer'].astype('float32')\n\n        # Giá trung bình mỗi article (popularity-weighted)\n        avg_price_art = (\n            df_master.groupby('article_id')['price']\n            .mean()\n            .reset_index()\n            .rename(columns={'price': 'avg_price_per_article'})\n        )\n        df_master = df_master.merge(avg_price_art, on='article_id', how='left')\n        df_master['avg_price_per_article'] = \\\n            df_master['avg_price_per_article'].astype('float32')\n\n        # Price ratio: giá giao dịch / avg giá khách hàng\n        df_master['price_ratio'] = (\n            df_master['price'] / df_master['avg_price_per_customer']\n        ).astype('float32')\n\n        # Price segment (categorical) — dựa trên phân vị toàn dataset\n        p33 = float(df_master['price'].quantile(0.33))\n        p66 = float(df_master['price'].quantile(0.66))\n\n        df_master['price_segment'] = 'premium'\n        df_master['price_segment'] = df_master['price_segment'].where(\n            df_master['price'] > p66, 'mid_range'\n        )\n        df_master['price_segment'] = df_master['price_segment'].where(\n            df_master['price'] > p33, 'budget'\n        )\n\n        # Customer price tier (theo avg của họ)\n        df_master['customer_price_tier'] = 'mid'\n        df_master['customer_price_tier'] = df_master['customer_price_tier'].where(\n            df_master['avg_price_per_customer'] > p66, 'premium'\n        )\n        df_master['customer_price_tier'] = df_master['customer_price_tier'].where(\n            df_master['avg_price_per_customer'] <= p33, 'budget'\n        )\n\n    else:\n        import polars as pl\n\n        # Avg price per customer\n        avg_cust = (\n            df_master.group_by('customer_id')\n            .agg(pl.col('price').mean().cast(pl.Float32).alias('avg_price_per_customer'))\n        )\n        df_master = df_master.join(avg_cust, on='customer_id', how='left')\n\n        # Avg price per article\n        avg_art = (\n            df_master.group_by('article_id')\n            .agg(pl.col('price').mean().cast(pl.Float32).alias('avg_price_per_article'))\n        )\n        df_master = df_master.join(avg_art, on='article_id', how='left')\n\n        # Price ratio\n        df_master = df_master.with_columns(\n            (pl.col('price') / pl.col('avg_price_per_customer'))\n            .cast(pl.Float32).alias('price_ratio')\n        )\n\n        # Price segments\n        p33 = df_master['price'].quantile(0.33)\n        p66 = df_master['price'].quantile(0.66)\n\n        df_master = df_master.with_columns(\n            pl.when(pl.col('price') <= p33).then(pl.lit('budget'))\n              .when(pl.col('price') <= p66).then(pl.lit('mid_range'))\n              .otherwise(pl.lit('premium'))\n              .alias('price_segment')\n        )\n\n        df_master = df_master.with_columns(\n            pl.when(pl.col('avg_price_per_customer') <= p33).then(pl.lit('budget'))\n              .when(pl.col('avg_price_per_customer') <= p66).then(pl.lit('mid'))\n              .otherwise(pl.lit('premium'))\n              .alias('customer_price_tier')\n        )\n\n    print(f\"  Thresholds: budget ≤ {p33:.4f} | mid ≤ {p66:.4f} | premium > {p66:.4f}\")\n    print(f\"\\n  Price segment distribution:\")\n    if use_gpu:\n        print(df_master['price_segment'].value_counts())\n    else:\n        print(\n            df_master.group_by('price_segment')\n            .agg(pl.count().alias('n'))\n            .sort('n', descending=True)\n        )\n\n    print(f\"\\n  Customer price tier distribution:\")\n    if use_gpu:\n        print(df_master['customer_price_tier'].value_counts())\n    else:\n        print(\n            df_master.group_by('customer_price_tier')\n            .agg(pl.count().alias('n'))\n            .sort('n', descending=True)\n        )\n\n    print(f\"\\n  New: avg_price_per_customer, avg_price_per_article,\")\n    print(f\"       price_ratio, price_segment, customer_price_tier\")\n\n    return df_master\n\ndf_master = add_price_features(df_master)\ngc.collect()\nprint(\"\\nPrice features done\")","metadata":{"trusted":true,"execution":{"iopub.status.busy":"2026-04-06T16:25:49.045343Z","iopub.execute_input":"2026-04-06T16:25:49.045969Z","iopub.status.idle":"2026-04-06T16:25:54.340940Z","shell.execute_reply.started":"2026-04-06T16:25:49.045938Z","shell.execute_reply":"2026-04-06T16:25:54.340178Z"}},"outputs":[],"execution_count":null},{"cell_type":"code","source":"# ============================================================\n# CELL 16: BƯỚC 5C — COLOUR & STYLE PREFERENCE FEATURES\n# ============================================================\nimport gc\n\ndef add_colour_style_features(df_master, use_gpu=USE_GPU):\n    print(\"=\" * 60)\n    print(\"BƯỚC 5C: COLOUR & STYLE FEATURES\")\n    print(\"=\" * 60)\n\n    if use_gpu:\n        import cudf\n\n        def customer_mode(col_name, alias):\n            \"\"\"Tìm giá trị xuất hiện nhiều nhất của col_name theo customer.\"\"\"\n            counts = (\n                df_master.groupby(['customer_id', col_name])\n                .size()\n                .reset_index(name='cnt')\n            )\n            idx = counts.groupby('customer_id')['cnt'].idxmax()\n            mode_df = counts.loc[idx][['customer_id', col_name]].rename(\n                columns={col_name: alias}\n            )\n            # Dọn rác ngay lập tức\n            del counts, idx\n            gc.collect()\n            return mode_df\n\n        # 1. Tính toán từng feature riêng lẻ\n        fav_colour  = customer_mode('perceived_colour_master_name', 'fav_colour')\n        fav_dept    = customer_mode('department_name',              'fav_department')\n        fav_group   = customer_mode('product_group_name',           'fav_product_group')\n        fav_garment = customer_mode('garment_group_name',           'fav_garment_group')\n\n        colour_div = (\n            df_master.groupby('customer_id')['perceived_colour_master_name']\n            .nunique()\n            .reset_index()\n            .rename(columns={'perceived_colour_master_name': 'colour_diversity'})\n        )\n        colour_div['colour_diversity'] = colour_div['colour_diversity'].astype('int8')\n\n        # 2. CHIẾN THUẬT: Gom tất cả vào 1 bảng profile nhỏ (~1.3 triệu dòng) TRƯỚC\n        cust_profile = fav_colour.merge(fav_dept, on='customer_id', how='outer')\n        cust_profile = cust_profile.merge(fav_group, on='customer_id', how='outer')\n        cust_profile = cust_profile.merge(fav_garment, on='customer_id', how='outer')\n        cust_profile = cust_profile.merge(colour_div, on='customer_id', how='outer')\n\n        # Ép kiểu category cho bảng profile để nhẹ nhất có thể\n        for col in ['fav_colour', 'fav_department', 'fav_product_group', 'fav_garment_group']:\n            cust_profile[col] = cust_profile[col].astype('category')\n\n        # Dọn rác các bảng con\n        del fav_colour, fav_dept, fav_group, fav_garment, colour_div\n        gc.collect()\n\n        # 3. CHỈ MERGE 1 LẦN DUY NHẤT vào df_master (31.7 triệu dòng)\n        df_master = df_master.merge(cust_profile, on='customer_id', how='left')\n        \n        del cust_profile\n        gc.collect()\n\n    else:\n        import polars as pl\n\n        def customer_mode_polars(df, col_name, alias):\n            return (\n                df.group_by(['customer_id', col_name])\n                .agg(pl.count().alias('cnt'))\n                .sort('cnt', descending=True)\n                .group_by('customer_id')\n                .agg(pl.first(col_name).alias(alias))\n            )\n\n        fav_colour  = customer_mode_polars(df_master, 'perceived_colour_master_name', 'fav_colour')\n        fav_dept    = customer_mode_polars(df_master, 'department_name', 'fav_department')\n        fav_group   = customer_mode_polars(df_master, 'product_group_name', 'fav_product_group')\n        fav_garment = customer_mode_polars(df_master, 'garment_group_name', 'fav_garment_group')\n\n        colour_div = (\n            df_master.group_by('customer_id')\n            .agg(pl.col('perceived_colour_master_name').n_unique()\n                 .cast(pl.Int8).alias('colour_diversity'))\n        )\n\n        # Polars cũng nên gom nhóm trước\n        cust_profile = fav_colour.join(fav_dept, on='customer_id', how='left')\n        cust_profile = cust_profile.join(fav_group, on='customer_id', how='left')\n        cust_profile = cust_profile.join(fav_garment, on='customer_id', how='left')\n        cust_profile = cust_profile.join(colour_div, on='customer_id', how='left')\n\n        df_master = df_master.join(cust_profile, on='customer_id', how='left')\n        \n        del fav_colour, fav_dept, fav_group, fav_garment, colour_div, cust_profile\n        gc.collect()\n\n    # Report\n    print(f\"  New customer-level features:\")\n    print(f\"    fav_colour, fav_department, fav_product_group,\")\n    print(f\"    fav_garment_group, colour_diversity\")\n\n    print(f\"\\n  Top 10 favourite colours (most common across customers):\")\n    if use_gpu:\n        print(df_master['fav_colour'].value_counts().head(10))\n    else:\n        print(\n            df_master.group_by('fav_colour')\n            .agg(pl.count().alias('n'))\n            .sort('n', descending=True)\n            .head(10)\n        )\n\n    print(f\"\\n  Colour diversity distribution:\")\n    if use_gpu:\n        print(df_master['colour_diversity'].describe())\n    else:\n        print(df_master['colour_diversity'].describe())\n\n    return df_master\n\ndf_master = add_colour_style_features(df_master)\ngc.collect()\nprint(\"\\nColour & style features done\")","metadata":{"trusted":true,"execution":{"iopub.status.busy":"2026-04-06T16:25:54.341960Z","iopub.execute_input":"2026-04-06T16:25:54.342249Z","iopub.status.idle":"2026-04-06T16:26:00.393454Z","shell.execute_reply.started":"2026-04-06T16:25:54.342224Z","shell.execute_reply":"2026-04-06T16:26:00.392824Z"}},"outputs":[],"execution_count":null},{"cell_type":"code","source":"# ============================================================\n# CELL 17: BƯỚC 5D — BEHAVIOUR FEATURES (RFM-inspired)\n# ============================================================\nimport gc\nimport numpy as np\n\ndef add_behaviour_features(df_master, use_gpu=USE_GPU):\n    print(\"=\" * 60)\n    print(\"BƯỚC 5D: BEHAVIOUR FEATURES (RFM)\")\n    print(\"=\" * 60)\n\n    if use_gpu:\n        import cudf\n\n        # 1. TẠO BẢNG RFM RIÊNG (Chỉ ~1.3 triệu dòng, thay vì 31.7 triệu dòng)\n        rfm = (\n            df_master.groupby('customer_id')\n            .agg({\n                'days_since_start': 'max',\n                'article_id': 'nunique',\n                'price': 'sum'\n            })\n            .reset_index()\n            .rename(columns={\n                'days_since_start': 'last_purchase_day',\n                'article_id': 'n_unique_articles',\n                'price': 'total_spend'\n            })\n        )\n\n        max_day = int(df_master['days_since_start'].max())\n        rfm['recency_days'] = (max_day - rfm['last_purchase_day']).astype('int16')\n\n        # Dtypes cho rfm\n        rfm['total_spend']       = rfm['total_spend'].astype('float32')\n        rfm['n_unique_articles'] = rfm['n_unique_articles'].astype('int16')\n        rfm['last_purchase_day'] = rfm['last_purchase_day'].astype('int16')\n\n        # 2. HÀM TÍNH ĐIỂM QUINTILE (Chỉ áp dụng trên 1.3 triệu dòng -> Rất an toàn cho RAM)\n        def quintile_score(series, reverse=False):\n            vals = series.to_pandas().fillna(0) # 1.3 triệu dòng nên kéo về Pandas vô tư\n            thresholds = np.percentile(vals, [20, 40, 60, 80])\n            if reverse:\n                thresholds = thresholds[::-1]\n            return thresholds\n\n        r_thresh = quintile_score(rfm['recency_days'],      reverse=True)\n        f_thresh = quintile_score(rfm['n_unique_articles'], reverse=False)\n        m_thresh = quintile_score(rfm['total_spend'],       reverse=False)\n\n        def make_score_col(col, thresholds, reverse=False):\n            scores = rfm[col].to_pandas().fillna(0)\n            if reverse:\n                result = np.where(scores <= thresholds[0], 5,\n                         np.where(scores <= thresholds[1], 4,\n                         np.where(scores <= thresholds[2], 3,\n                         np.where(scores <= thresholds[3], 2, 1))))\n            else:\n                result = np.where(scores <= thresholds[0], 1,\n                         np.where(scores <= thresholds[1], 2,\n                         np.where(scores <= thresholds[2], 3,\n                         np.where(scores <= thresholds[3], 4, 5))))\n            return cudf.Series(result.astype('int8'))\n\n        # 3. GÁN ĐIỂM RFM LÊN BẢNG NHỎ\n        rfm['r_score'] = make_score_col('recency_days',      r_thresh, reverse=True)\n        rfm['f_score'] = make_score_col('n_unique_articles', f_thresh)\n        rfm['m_score'] = make_score_col('total_spend',       m_thresh)\n        rfm['rfm_score'] = (\n            rfm['r_score'].astype('int16') +\n            rfm['f_score'].astype('int16') +\n            rfm['m_score'].astype('int16')\n        ).astype('int8')\n\n        # 4. LẤY TOP 10K TỪ BẢNG NHỎ\n        top10k = rfm.nlargest(10_000, 'rfm_score')[['customer_id', 'rfm_score', 'recency_days', 'total_spend']]\n\n        # 5. MERGE 1 LẦN DUY NHẤT VÀO DF_MASTER\n        df_master = df_master.merge(rfm, on='customer_id', how='left')\n\n        # Dọn rác siêu kỹ\n        del rfm\n        gc.collect()\n\n    else:\n        import polars as pl\n\n        # Tuơng tự cho CPU, gom nhóm trước khi tính điểm\n        rfm = (\n            df_master.group_by('customer_id')\n            .agg([\n                pl.col('days_since_start').max().alias('last_purchase_day'),\n                pl.col('article_id').n_unique().cast(pl.Int16).alias('n_unique_articles'),\n                pl.col('price').sum().cast(pl.Float32).alias('total_spend'),\n            ])\n        )\n        \n        max_day = df_master['days_since_start'].max()\n        rfm = rfm.with_columns(\n            (max_day - pl.col('last_purchase_day')).cast(pl.Int16).alias('recency_days')\n        )\n\n        def make_score_expr(col, thresholds, reverse=False):\n            t = thresholds\n            if reverse:\n                return (\n                    pl.when(pl.col(col) <= t[0]).then(5)\n                      .when(pl.col(col) <= t[1]).then(4)\n                      .when(pl.col(col) <= t[2]).then(3)\n                      .when(pl.col(col) <= t[3]).then(2)\n                      .otherwise(1).cast(pl.Int8)\n                )\n            else:\n                return (\n                    pl.when(pl.col(col) <= t[0]).then(1)\n                      .when(pl.col(col) <= t[1]).then(2)\n                      .when(pl.col(col) <= t[2]).then(3)\n                      .when(pl.col(col) <= t[3]).then(4)\n                      .otherwise(5).cast(pl.Int8)\n                )\n\n        # Cần lấy array để tính percentile\n        r_thresh = np.percentile(rfm['recency_days'].to_numpy(), [20, 40, 60, 80])[::-1]\n        f_thresh = np.percentile(rfm['n_unique_articles'].to_numpy(), [20, 40, 60, 80])\n        m_thresh = np.percentile(rfm['total_spend'].to_numpy(), [20, 40, 60, 80])\n\n        rfm = rfm.with_columns([\n            make_score_expr('recency_days',      r_thresh, reverse=True).alias('r_score'),\n            make_score_expr('n_unique_articles', f_thresh).alias('f_score'),\n            make_score_expr('total_spend',       m_thresh).alias('m_score'),\n        ])\n        \n        rfm = rfm.with_columns(\n            (pl.col('r_score').cast(pl.Int16) +\n             pl.col('f_score').cast(pl.Int16) +\n             pl.col('m_score').cast(pl.Int16)).cast(pl.Int8).alias('rfm_score')\n        )\n\n        top10k = rfm.sort('rfm_score', descending=True).head(10_000).select(['customer_id', 'rfm_score', 'recency_days', 'total_spend'])\n        df_master = df_master.join(rfm, on='customer_id', how='left')\n        \n        del rfm\n        gc.collect()\n\n    print(f\"  New features: last_purchase_day, n_unique_articles, total_spend,\")\n    print(f\"                recency_days, r_score, f_score, m_score, rfm_score\")\n\n    # Lưu Parquet (Đã kèm thuốc trị lỗi category)\n    top10k_path = f'{OUTPUT_DIR}/top10k_active_users.parquet'\n    if use_gpu:\n        cat_cols = top10k.select_dtypes(include=['category']).columns\n        for col in cat_cols:\n            top10k[col] = top10k[col].astype(str)\n        top10k.to_parquet(top10k_path)\n    else:\n        top10k.write_parquet(top10k_path)\n\n    print(f\"\\n  Top-10K active users saved → {top10k_path}\")\n\n    return df_master\n\ndf_master = add_behaviour_features(df_master)\ngc.collect()\nprint(\"\\nBehaviour features (RFM) done\")","metadata":{"trusted":true,"execution":{"iopub.status.busy":"2026-04-06T16:26:00.394356Z","iopub.execute_input":"2026-04-06T16:26:00.394641Z","iopub.status.idle":"2026-04-06T16:26:03.289469Z","shell.execute_reply.started":"2026-04-06T16:26:00.394613Z","shell.execute_reply":"2026-04-06T16:26:03.288756Z"}},"outputs":[],"execution_count":null},{"cell_type":"code","source":"# ============================================================\n# CELL 18: SAVE df_features — checkpoint cuối Chương 2\n# ============================================================\nimport os\nimport gc\n\ndef save_features_checkpoint(df_master, use_gpu=USE_GPU):\n    path = f'{OUTPUT_DIR}/df_features.parquet'\n\n    if use_gpu:\n        import pyarrow.parquet as pq\n        print(\"  [GPU Mode] Chuyển đổi qua PyArrow để lưu Category an toàn...\")\n        \n        # Chuyển cuDF DataFrame sang PyArrow Table (chuyển memory sang CPU dạng nén)\n        # PyArrow giữ nguyên được định dạng category mà KHÔNG làm phình RAM\n        arrow_table = df_master.to_arrow()\n        \n        print(\"  [GPU Mode] Đang ghi file Parquet...\")\n        pq.write_table(arrow_table, path)\n        \n        # Dọn rác\n        del arrow_table\n        gc.collect()\n    else:\n        # Polars CPU tự hỗ trợ rất tốt rồi\n        df_master.write_parquet(path)\n\n    size_mb = os.path.getsize(path) / 1e6\n\n    print(\"=\" * 60)\n    print(\"SAVE: df_features.parquet\")\n    print(\"=\" * 60)\n    print(f\"  Path  : {path}\")\n    print(f\"  Size  : {size_mb:.1f} MB\")\n    print(f\"  Shape : {df_master.shape}\")\n\n    print(f\"\\n  All columns ({df_master.shape[1]}):\")\n    for i, col in enumerate(df_master.columns, 1):\n        print(f\"    {i:3d}. {col}\")\n\nsave_features_checkpoint(df_master)","metadata":{"trusted":true,"execution":{"iopub.status.busy":"2026-04-06T16:26:03.290543Z","iopub.execute_input":"2026-04-06T16:26:03.290847Z","iopub.status.idle":"2026-04-06T16:27:19.670772Z","shell.execute_reply.started":"2026-04-06T16:26:03.290825Z","shell.execute_reply":"2026-04-06T16:27:19.670032Z"}},"outputs":[],"execution_count":null},{"cell_type":"code","source":"# ============================================================\n# CELL 19: BƯỚC 6 — CRITICAL THINKING SUMMARY\n# Tiền xử lý Chương 2 phục vụ gì cho Chương 3 & 4?\n# ============================================================\n\ncritical_thinking_summary = \"\"\"\n╔══════════════════════════════════════════════════════════════╗\n║       BƯỚC 6: CRITICAL THINKING — KẾT NỐI CHƯƠNG 2→3→4     ║\n╚══════════════════════════════════════════════════════════════╝\n\n━━━ A. CHƯƠNG 2 → CHƯƠNG 3 (Lambda Architecture) ━━━━━━━━━━━━\n\n  Feature / Quyết định          │ Phục vụ gì ở Chương 3?\n  ──────────────────────────────┼────────────────────────────────────\n  top10k_active_users (RFM)    │ Chọn đúng 10K users để lưu\n                                │ recommendation lên MongoDB (<500MB)\n  Downcast int8/float32         │ Parquet nhỏ → ingestion Kafka nhanh\n  buyer_type flag               │ Batch layer lọc reseller trước khi\n                                │ tính association rules\n  df_features.parquet           │ Input cho Spark batch job\n  fashion_season / quarter      │ Partition key cho Spark để xử lý\n                                │ theo mùa (reduce shuffle)\n\n━━━ B. CHƯƠNG 2 → CHƯƠNG 4 (Đánh giá mô hình) ━━━━━━━━━━━━━\n\n  Feature / Quyết định          │ Phục vụ gì ở Chương 4?\n  ──────────────────────────────┼────────────────────────────────────\n  days_since_start              │ Time-based train/test split —\n                                │ train = mọi ngày trừ 7 cuối\n  week_of_year                  │ Detect seasonal patterns trong\n                                │ association rules (rules mùa SS ≠ FW)\n  age_group, price_tier         │ Stratified evaluation: MAP@12 theo\n                                │ từng segment (teen vs. adult)\n  fav_colour, fav_department    │ Feature cho candidate generation\n                                │ (lọc items phù hợp style trước)\n  rfm_score                     │ Weight transactions: loyal user\n                                │ signal mạnh hơn → higher weight\n                                │ trong support calculation\n  buyer_type='normal'           │ Chỉ dùng normal buyers để train\n                                │ FP-Growth (tránh reseller bias)\n  colour_diversity              │ Personalization depth metric:\n                                │ KH narrow taste → recommend similar\n                                │ KH broad taste → recommend diverse\n\n━━━ C. QUYẾT ĐỊNH THIẾT KẾ QUAN TRỌNG NHẤT ━━━━━━━━━━━━━━━━\n\n  1. KHÔNG DROP null → GIỮ TOÀN BỘ 31M transactions\n     Lý do: Trong retail, mỗi transaction là ground truth.\n     Drop = mất signal. Null được encode thành trạng thái riêng.\n\n  2. LEFT JOIN (không INNER) ở bước 2\n     Lý do: Tránh mất transaction của articles/customers có\n     data lag. Trong production, catalog thay đổi liên tục.\n\n  3. Downcast TRƯỚC feature engineering\n     Lý do: Tránh OOM khi tạo các cột tổng hợp (avg, nunique)\n     trên 31M dòng với float64 — dễ vượt 16GB VRAM T4.\n\n  4. RFM score → chọn 10K users cho MongoDB\n     Lý do: 10K × 12 recommendations × ~100 bytes/doc ≈ 12MB\n     Nằm trong ngưỡng 500MB MongoDB Atlas free tier.\n     RFM đảm bảo chọn đúng users có giá trị kinh doanh cao nhất.\n\n  5. buyer_type flag (không drop outlier)\n     Lý do: Reseller pattern là noise với recommendation nhưng\n     là signal với fraud detection & inventory planning.\n     Drop = mất thông tin quý cho downstream analysis.\n\"\"\"\n\nprint(critical_thinking_summary)","metadata":{"trusted":true,"execution":{"iopub.status.busy":"2026-04-06T16:27:19.671736Z","iopub.execute_input":"2026-04-06T16:27:19.671994Z","iopub.status.idle":"2026-04-06T16:27:19.678053Z","shell.execute_reply.started":"2026-04-06T16:27:19.671971Z","shell.execute_reply":"2026-04-06T16:27:19.677247Z"}},"outputs":[],"execution_count":null},{"cell_type":"code","source":"# ============================================================\n# CELL 20: VISUALISATION — Feature engineering overview chart\n# ============================================================\nimport matplotlib.pyplot as plt\nimport matplotlib.patches as mpatches\nimport numpy as np\nimport polars as pl # Thêm import dự phòng nếu chạy CPU\n\nfig, axes = plt.subplots(1, 3, figsize=(18, 6))\nfig.suptitle('Chương 2 — Feature Engineering Summary', fontsize=14, fontweight='bold')\n\n# ── Plot 1: RFM Score Distribution ──\nax = axes[0]\nif USE_GPU:\n    # Bỏ dòng rfm_vals bị lỗi, dùng trực tiếp drop_duplicates cho DataFrame\n    rfm_pd = df_master[['customer_id','rfm_score']].drop_duplicates(subset=['customer_id'])\n    rfm_pd = rfm_pd['rfm_score'].to_pandas()\nelse:\n    rfm_pd = (\n        df_master.select(['customer_id','rfm_score'])\n        .unique('customer_id')['rfm_score']\n        .to_numpy()\n    )\n\nax.hist(rfm_pd, bins=13, range=(3,16), color='#3266ad', alpha=0.85,\n        edgecolor='white', linewidth=0.5)\nax.set_title('RFM Score Distribution\\n(per customer)')\nax.set_xlabel('RFM Score (3=worst, 15=best)')\nax.set_ylabel('Số khách hàng')\nax.axvline(np.median(rfm_pd), color='red', linestyle='--', linewidth=1.5,\n           label=f'Median: {np.median(rfm_pd):.0f}')\nax.legend(fontsize=9)\n\n# ── Plot 2: Fashion Season vs Transaction Count ──\nax = axes[1]\nif USE_GPU:\n    season_counts = df_master['fashion_season'].value_counts().to_pandas()\n    # FIX: value_counts() trả về Series, phải lấy từ index và values\n    s_labels = season_counts.index.tolist()\n    s_values = season_counts.values.tolist()\nelse:\n    season_counts = (\n        df_master.group_by('fashion_season')\n        .agg(pl.count().alias('n'))\n        .sort('n', descending=True)\n        .to_pandas()\n    )\n    s_labels = season_counts['fashion_season'].tolist()\n    s_values = season_counts['n'].tolist()\n\ncolors_season = ['#e67e22' if s == 'FW' else '#3498db' for s in s_labels]\nbars = ax.bar(s_labels, s_values, color=colors_season, alpha=0.85,\n              edgecolor='white', linewidth=0.5)\nax.set_title('Transactions by\\nFashion Season')\nax.set_ylabel('Số giao dịch')\nimport matplotlib.ticker as mticker\nax.yaxis.set_major_formatter(mticker.FuncFormatter(lambda x, _: f'{x/1e6:.1f}M'))\npatches = [\n    mpatches.Patch(color='#e67e22', label='FW: Fall/Winter (Aug–Jan)'),\n    mpatches.Patch(color='#3498db', label='SS: Spring/Summer (Feb–Jul)'),\n]\nax.legend(handles=patches, fontsize=8)\n\n# ── Plot 3: Customer Price Tier Breakdown ──\nax = axes[2]\nif USE_GPU:\n    tier_df = (\n        df_master[['customer_id', 'customer_price_tier']]\n        .drop_duplicates(subset=['customer_id'])['customer_price_tier']\n        .value_counts()\n        .to_pandas()\n    )\n    # FIX: value_counts() trả về Series, phải lấy từ index và values\n    t_labels = tier_df.index.tolist()\n    t_values = tier_df.values.tolist()\nelse:\n    tier_df = (\n        df_master.select(['customer_id', 'customer_price_tier'])\n        .unique('customer_id')\n        .group_by('customer_price_tier')\n        .agg(pl.count().alias('n'))\n        .sort('n', descending=True)\n        .to_pandas()\n    )\n    t_labels = tier_df['customer_price_tier'].tolist()\n    t_values = tier_df['n'].tolist()\n\ntier_colors = {'budget': '#27ae60', 'mid': '#f39c12', 'premium': '#8e44ad'}\nbar_colors  = [tier_colors.get(l, '#7f8c8d') for l in t_labels]\nax.barh(t_labels, t_values, color=bar_colors, alpha=0.85,\n        edgecolor='white', linewidth=0.5)\nax.set_title('Customer Price Tier\\n(unique customers)')\nax.set_xlabel('Số khách hàng')\nax.xaxis.set_major_formatter(mticker.FuncFormatter(lambda x, _: f'{x/1e3:.0f}K'))\n\nplt.tight_layout()\nplt.savefig(f'{OUTPUT_DIR}/feature_engineering_summary.png', dpi=150, bbox_inches='tight')\nplt.show()\nprint(\"Feature engineering chart saved!\")","metadata":{"trusted":true,"execution":{"iopub.status.busy":"2026-04-06T16:27:19.678879Z","iopub.execute_input":"2026-04-06T16:27:19.679148Z","iopub.status.idle":"2026-04-06T16:27:21.026142Z","shell.execute_reply.started":"2026-04-06T16:27:19.679116Z","shell.execute_reply":"2026-04-06T16:27:21.025341Z"}},"outputs":[],"execution_count":null},{"cell_type":"markdown","source":"# **CHƯƠNG 2: Lambda Architecture: Kafka -> Spark -> Mongo DB**","metadata":{}},{"cell_type":"code","source":"# ============================================================\n# CELL 21: CÀI THƯ VIỆN & CẤU HÌNH KẾT NỐI\n# ============================================================\nimport subprocess, sys\n\n# Cài các thư viện cần thiết\npkgs = ['pymongo[srv]', 'kafka-python', 'pyspark', 'dnspython']\nfor pkg in pkgs:\n    subprocess.run([sys.executable, '-m', 'pip', 'install', pkg, '-q'],\n                   capture_output=True)\n\nprint(\"Libraries installed\")\n\n# ── Cấu hình MongoDB Atlas ──\n# Thay MONGO_URI bằng connection string thực của bạn\n# Format: mongodb+srv://<user>:<password>@<cluster>.mongodb.net/\n\nMONGO_URI = \"mongodb+srv://phuonganh0805:la080503@clusterbigdata85.nyncx70.mongodb.net/?appName=ClusterBigData85\"\nDB_NAME   = \"hm_recommendation\"\n\n# Collections (chỉ lưu aggregated results — không bao giờ lưu raw data)\nCOL_RULES  = \"association_rules\"   # FP-Growth rules\nCOL_RECS   = \"user_recommendations\" # Top-12 per user\nCOL_TRENDS = \"item_trends\"          # Popularity theo thời gian\nCOL_META   = \"system_metadata\"      # Audit log\n\n# ── Kết nối MongoDB ──\nfrom pymongo import MongoClient\nfrom pymongo.errors import ConnectionFailure\nimport certifi\n\ndef get_mongo_client():\n    \"\"\"\n    BIỆN LUẬN: Dùng connection pooling (default của pymongo)\n    thay vì tạo connection mới mỗi lần. Trong serving layer,\n    mỗi query recommendation cần latency <100ms — connection\n    overhead sẽ giết chết SLA nếu không pool.\n    \"\"\"\n    client = MongoClient(\n        MONGO_URI,\n        tlsCAFile=certifi.where(),\n        serverSelectionTimeoutMS=5000,\n        maxPoolSize=10,       # giới hạn connection pool\n        connectTimeoutMS=3000,\n    )\n    try:\n        client.admin.command('ping')\n        print(\"MongoDB Atlas connected!\")\n    except ConnectionFailure as e:\n        print(f\"Connection failed: {e}\")\n        raise\n    return client\n\nclient = get_mongo_client()\ndb     = client[DB_NAME]\n\n# Tạo indexes ngay khi khởi tạo (quan trọng cho query speed)\ndef create_indexes():\n    \"\"\"\n    BIỆN LUẬN: Index trên customer_id giúp lookup recommendation\n    O(log n) thay vì O(n). Với 10K users, không critical, nhưng\n    đây là best practice cho production-grade system.\n    \"\"\"\n    # user_recommendations: lookup theo customer_id\n    db[COL_RECS].create_index(\"customer_id\", unique=True)\n\n    # association_rules: query theo antecedent\n    db[COL_RULES].create_index([(\"antecedents\", 1)])\n    db[COL_RULES].create_index([(\"lift\", -1)])        # sort by lift\n\n    # item_trends: lookup theo article_id + week\n    db[COL_TRENDS].create_index([(\"article_id\", 1), (\"week\", 1)])\n\n    print(\"Indexes created\")\n\ncreate_indexes()","metadata":{"trusted":true,"execution":{"iopub.status.busy":"2026-04-06T16:28:40.771217Z","iopub.execute_input":"2026-04-06T16:28:40.771535Z","iopub.status.idle":"2026-04-06T16:28:59.189629Z","shell.execute_reply.started":"2026-04-06T16:28:40.771507Z","shell.execute_reply":"2026-04-06T16:28:59.189003Z"}},"outputs":[],"execution_count":null},{"cell_type":"code","source":"# ============================================================\n# CELL 22: INGESTION LAYER — Kafka Mock\n#\n# BIỆN LUẬN: Kafka là industry standard cho event streaming\n# trong retail (Zalando, H&M thực tế đều dùng Kafka).\n# Trên Kaggle không có Kafka server thật, ta mock bằng\n# in-memory queue — logic producer/consumer hoàn toàn giữ\n# nguyên, chỉ thay transport layer. Điều này đủ để minh chứng\n# kiến trúc mà không cần infra phức tạp.\n# ============================================================\nimport json\nimport time\nimport threading\nimport queue\nfrom datetime import datetime\nimport random\n\n# ── Mock Kafka: dùng Python queue thay broker ──\nkafka_topic = queue.Queue(maxsize=10_000)\n\nclass MockKafkaProducer:\n    \"\"\"Giả lập Kafka producer: gửi purchase events vào topic.\"\"\"\n\n    def __init__(self, topic_queue: queue.Queue):\n        self.topic   = topic_queue\n        self.sent    = 0\n        self.dropped = 0\n\n    def send(self, event: dict):\n        msg = {\n            'timestamp': datetime.utcnow().isoformat(),\n            'event_type': 'purchase',\n            'payload': event,\n        }\n        try:\n            self.topic.put_nowait(json.dumps(msg))\n            self.sent += 1\n        except queue.Full:\n            self.dropped += 1   # backpressure — topic đầy\n\n    def flush(self):\n        print(f\"  Producer: sent={self.sent:,}, dropped={self.dropped:,}\")\n\n\nclass MockKafkaConsumer:\n    \"\"\"Giả lập Kafka consumer: đọc events từ topic, batch xử lý.\"\"\"\n\n    def __init__(self, topic_queue: queue.Queue, batch_size=500):\n        self.topic      = topic_queue\n        self.batch_size = batch_size\n        self.batches    = []\n        self.total      = 0\n\n    def poll_batch(self, timeout=2.0):\n        \"\"\"Đọc tối đa batch_size messages từ topic.\"\"\"\n        batch  = []\n        t_stop = time.time() + timeout\n        while len(batch) < self.batch_size and time.time() < t_stop:\n            try:\n                raw = self.topic.get(timeout=0.1)\n                msg = json.loads(raw)\n                batch.append(msg['payload'])\n            except queue.Empty:\n                break\n        if batch:\n            self.batches.append(batch)\n            self.total += len(batch)\n        return batch\n\n\ndef simulate_kafka_ingestion(df_master, n_events=50_000, use_gpu=USE_GPU):\n    \"\"\"\n    Mô phỏng luồng purchase events từ H&M transactions.\n    Producer gửi N events, consumer nhận theo batch.\n    \"\"\"\n    print(\"=\" * 60)\n    print(\"INGESTION LAYER: Kafka simulation\")\n    print(\"=\" * 60)\n    print(f\"  Simulating {n_events:,} purchase events...\")\n\n    producer = MockKafkaProducer(kafka_topic)\n    consumer = MockKafkaConsumer(kafka_topic, batch_size=1000)\n\n    # Sample transactions để stream\n    if use_gpu:\n        sample = df_master.sample(n=min(n_events, len(df_master)),\n                                  random_state=42).to_pandas()\n    else:\n        sample = df_master.sample(n=min(n_events, len(df_master)),\n                                  seed=42).to_pandas()\n\n    t0 = time.time()\n\n    # Producer thread: gửi events\n    def produce():\n        for _, row in sample.iterrows():\n            event = {\n                'customer_id' : str(row['customer_id']),\n                'article_id'  : int(row['article_id']),\n                'price'       : float(row.get('price', 0)),\n                'channel'     : int(row.get('sales_channel_id', 1)),\n                'week'        : int(row.get('week_of_year', 1)),\n                'season'      : str(row.get('fashion_season', 'SS')),\n            }\n            producer.send(event)\n            # Throttle: simulate ~5K events/sec\n            if producer.sent % 5000 == 0:\n                time.sleep(0.001)\n        producer.flush()\n\n    # Consumer: đọc batch trong main thread\n    prod_thread = threading.Thread(target=produce, daemon=True)\n    prod_thread.start()\n\n    all_events = []\n    print(\"  Consuming batches\", end='', flush=True)\n    while prod_thread.is_alive() or not kafka_topic.empty():\n        batch = consumer.poll_batch(timeout=1.0)\n        if batch:\n            all_events.extend(batch)\n            print('.', end='', flush=True)\n\n    prod_thread.join()\n    elapsed = time.time() - t0\n\n    print(f\"\\n\\nIngestion complete\")\n    print(f\"  Total events consumed : {len(all_events):,}\")\n    print(f\"  Total batches         : {len(consumer.batches)}\")\n    print(f\"  Throughput            : {len(all_events)/elapsed:,.0f} events/sec\")\n    print(f\"  Elapsed               : {elapsed:.2f}s\")\n\n    # Sample event để kiểm tra\n    print(f\"\\n  Sample event:\")\n    print(json.dumps(all_events[0], indent=4))\n\n    return all_events\n\nkafka_events = simulate_kafka_ingestion(df_master, n_events=50_000)","metadata":{"trusted":true,"execution":{"iopub.status.busy":"2026-04-06T16:30:28.091319Z","iopub.execute_input":"2026-04-06T16:30:28.091881Z","iopub.status.idle":"2026-04-06T16:30:34.509805Z","shell.execute_reply.started":"2026-04-06T16:30:28.091852Z","shell.execute_reply":"2026-04-06T16:30:34.509177Z"}},"outputs":[],"execution_count":null},{"cell_type":"code","source":"# ============================================================\n# CELL 23: BATCH LAYER — PySpark FP-Growth\n#\n# BIỆN LUẬN: Spark MLlib FP-Growth xử lý parallel trên\n# cluster/GPU tốt hơn single-machine mlxtend khi dataset >10M.\n# Ta dùng Spark để giữ đúng tinh thần Big Data architecture:\n# batch job chạy định kỳ (daily), kết quả ghi vào serving layer.\n# Kaggle có thể chạy Spark local[*] tận dụng cả 2 CPU core T4.\n# ============================================================\nimport time\nfrom pyspark.sql import SparkSession\nfrom pyspark.sql import functions as F\nfrom pyspark.ml.fpm import FPGrowth\nimport pandas as pd\n\ndef create_spark_session():\n    spark = (\n        SparkSession.builder\n        .appName(\"HM_FPGrowth_Batch\")\n        .master(\"local[*]\")           # Dùng tất cả CPU cores\n        .config(\"spark.driver.memory\", \"8g\")\n        .config(\"spark.executor.memory\", \"8g\")\n        .config(\"spark.sql.shuffle.partitions\", \"50\")\n        .config(\"spark.default.parallelism\", \"50\")\n        # Tối ưu cho H&M dataset\n        .config(\"spark.sql.adaptive.enabled\", \"true\")\n        .config(\"spark.sql.adaptive.coalescePartitions.enabled\", \"true\")\n        .getOrCreate()\n    )\n    spark.sparkContext.setLogLevel(\"WARN\")\n    print(f\"Spark {spark.version} started — local[*]\")\n    return spark\n\nspark = create_spark_session()\n\ndef prepare_transactions_for_fpgrowth(df_master, use_gpu=USE_GPU,\n                                       min_basket_size=2,\n                                       sample_frac=0.3):\n    \"\"\"\n    Chuyển df_master → format basket (customer_id, [article_ids])\n    cho FP-Growth.\n    \"\"\"\n    print(\"  Preparing baskets for FP-Growth...\")\n\n    if use_gpu:\n        # FIX: Kiểm tra xem có cột buyer_type không trước khi lọc\n        if 'buyer_type' in df_master.columns:\n            df_normal = df_master[df_master['buyer_type'] == 'normal']\n        else:\n            df_normal = df_master\n            \n        df_pd = df_normal[['customer_id', 'article_id']].to_pandas()\n    else:\n        import polars as pl\n        if 'buyer_type' in df_master.columns:\n            df_pd = (\n                df_master\n                .filter(pl.col('buyer_type') == 'normal')\n                .select(['customer_id', 'article_id'])\n                .to_pandas()\n            )\n        else:\n            df_pd = df_master.select(['customer_id', 'article_id']).to_pandas()\n\n    # Sample để tăng tốc\n    df_pd = df_pd.sample(frac=sample_frac, random_state=42)\n\n    # Group thành baskets\n    baskets = (\n        df_pd.groupby('customer_id')['article_id']\n        .apply(lambda x: list(x.astype(str).unique()))\n        .reset_index()\n        .rename(columns={'article_id': 'items'})\n    )\n\n    # Lọc basket có ít nhất min_basket_size items\n    baskets = baskets[baskets['items'].apply(len) >= min_basket_size]\n\n    print(f\"  Total baskets : {len(baskets):,}\")\n    print(f\"  Avg basket sz : {baskets['items'].apply(len).mean():.1f}\")\n    print(f\"  Max basket sz : {baskets['items'].apply(len).max()}\")\n\n    # Chuyển sang Spark DataFrame\n    spark_df = spark.createDataFrame(baskets)\n    print(f\"  Spark partitions: {spark_df.rdd.getNumPartitions()}\")\n\n    return spark_df\n\nspark_baskets = prepare_transactions_for_fpgrowth(df_master)\n\n\ndef run_fpgrowth_batch(spark_df, min_support=0.01, min_confidence=0.3):\n    \"\"\"\n    Chạy FP-Growth trên Spark, trả về association rules DataFrame.\n    \"\"\"\n    print(f\"\\n  Running FP-Growth...\")\n    print(f\"  min_support={min_support}, min_confidence={min_confidence}\")\n\n    t0 = time.time()\n    fp = FPGrowth(\n        itemsCol=\"items\",\n        minSupport=min_support,\n        minConfidence=min_confidence,\n        numPartitions=50,\n    )\n    model     = fp.fit(spark_df)\n    rules_sdf = model.associationRules\n    freq_sdf  = model.freqItemsets\n\n    # Cache để tránh recompute\n    rules_sdf.cache()\n    freq_sdf.cache()\n\n    n_rules = rules_sdf.count()\n    n_freq  = freq_sdf.count()\n    elapsed = time.time() - t0\n\n    print(f\"  FP-Growth done in {elapsed:.1f}s\")\n    print(f\"  Frequent itemsets : {n_freq:,}\")\n    print(f\"  Association rules : {n_rules:,}\")\n\n    # Thêm Lift filter: chỉ giữ lift > 1 (positive association)\n    rules_filtered = rules_sdf.filter(F.col(\"lift\") > 1.0)\n    n_lift = rules_filtered.count()\n    print(f\"  Rules (lift>1)    : {n_lift:,}\")\n\n    # Top 10 rules by lift\n    print(f\"\\n  Top 10 rules by lift:\")\n    rules_filtered.orderBy(F.col(\"lift\").desc()).show(10, truncate=50)\n\n    return model, rules_filtered, freq_sdf\n\nspark_model, spark_rules, spark_freq = run_fpgrowth_batch(\n    spark_baskets, \n    min_support=0.0005,    # Giảm từ 1% xuống 0.05%\n    min_confidence=0.05    # Giảm từ 30% xuống 5%\n)","metadata":{"trusted":true,"execution":{"iopub.status.busy":"2026-04-06T16:41:37.786315Z","iopub.execute_input":"2026-04-06T16:41:37.786899Z","iopub.status.idle":"2026-04-06T16:52:38.549442Z","shell.execute_reply.started":"2026-04-06T16:41:37.786867Z","shell.execute_reply":"2026-04-06T16:52:38.548494Z"}},"outputs":[],"execution_count":null},{"cell_type":"code","source":"# ============================================================\n# CELL 24: SERVING LAYER — MongoDB Schema Design\n#\n# BIỆN LUẬN — Embedding (lồng nhau) thay vì Referencing:\n# Với giới hạn 500MB MongoDB Atlas free, ta dùng embedded\n# documents để tránh JOIN (MongoDB không có JOIN hiệu quả).\n# Top-12 recommendations lồng thẳng vào user document →\n# 1 query = lấy được toàn bộ info cần thiết cho API response.\n# Estimate: 10K users × 1KB/doc ≈ 10MB. Rules: ~5K × 200B ≈ 1MB.\n# Tổng << 500MB, an toàn.\n# ============================================================\nfrom pymongo import UpdateOne, InsertOne\nfrom datetime import datetime, timezone\nimport math\n\n# ── Schema definitions ──\n\nRULE_SCHEMA = {\n    # association_rules collection\n    \"antecedents\"  : [\"article_id_1\", \"article_id_2\"],  # list of strings\n    \"consequents\"  : [\"article_id_3\"],\n    \"support\"      : 0.012,      # float\n    \"confidence\"   : 0.45,       # float\n    \"lift\"         : 2.3,        # float\n    \"conviction\"   : 1.8,        # float (optional)\n    \"season\"       : \"SS\",       # tag mùa\n    \"created_at\"   : datetime.now(timezone.utc),\n    \"batch_run_id\" : \"2024-w52-batch\",\n}\n\nREC_SCHEMA = {\n    # user_recommendations collection\n    \"customer_id\"      : \"abc123def456\",\n    \"rfm_score\"        : 12,\n    \"customer_tier\"    : \"premium\",\n    \"top12_articles\"   : [\n        {\n            \"rank\"        : 1,\n            \"article_id\"  : \"0738413001\",\n            \"score\"       : 0.92,\n            \"method\"      : \"fp_growth\",   # hoặc \"popularity\", \"collab\"\n            \"season_match\": True,\n        },\n        # ... 11 more\n    ],\n    \"fav_colour\"       : \"Black\",\n    \"fav_department\"   : \"Ladieswear\",\n    \"last_updated\"     : datetime.now(timezone.utc),\n    \"week_generated\"   : 52,\n}\n\n\ndef write_association_rules_to_mongo(spark_rules, db, top_n=500):\n    \"\"\"\n    Ghi top N association rules lên MongoDB.\n\n    BIỆN LUẬN: Chỉ lưu top 500 rules (by lift) thay vì tất cả.\n    Serving layer chỉ cần rules chất lượng cao để lookup real-time.\n    Full rule set vẫn lưu trong Spark (batch layer) để audit.\n    \"\"\"\n    print(\"=\" * 60)\n    print(\"SERVING: Write association rules → MongoDB\")\n    print(\"=\" * 60)\n\n    # Chuyển Spark → Python list\n    rules_pd = (\n        spark_rules\n        .orderBy(F.col(\"lift\").desc())\n        .limit(top_n)\n        .toPandas()\n    )\n\n    print(f\"  Rules to write: {len(rules_pd)}\")\n\n    # Xây documents\n    docs = []\n    now  = datetime.now(timezone.utc)\n    for _, row in rules_pd.iterrows():\n        # Conviction = P(not C) / P(not C | A) = (1-conf)/(1-support)\n        try:\n            conviction = (1 - row['confidence']) / (1 - row['support'])\n        except ZeroDivisionError:\n            conviction = float('inf')\n\n        doc = {\n            # FIX LỖI Ở 2 DÒNG NÀY: Dùng antecedent và consequent (không có s)\n            \"antecedents\"  : list(row['antecedent']),\n            \"consequents\"  : list(row['consequent']),\n            \"support\"      : round(float(row['support']),    6),\n            \"confidence\"   : round(float(row['confidence']), 6),\n            \"lift\"         : round(float(row['lift']),       6),\n            \"conviction\"   : round(float(conviction),        6),\n            \"created_at\"   : now,\n            \"batch_run_id\" : f\"kaggle-{now.strftime('%Y%m%d')}\",\n        }\n        docs.append(doc)\n\n    # Bulk insert (xoá cũ trước để idempotent)\n    col = db[COL_RULES]\n    col.delete_many({\"batch_run_id\": docs[0][\"batch_run_id\"]})\n\n    result = col.insert_many(docs, ordered=False)\n    print(f\"  Inserted {len(result.inserted_ids)} rules\")\n\n    # Verify\n    sample = col.find_one({}, sort=[(\"lift\", -1)])\n    print(f\"\\n  Sample rule (highest lift):\")\n    print(f\"    {sample['antecedents']} → {sample['consequents']}\")\n    print(f\"    support={sample['support']:.4f}, confidence={sample['confidence']:.4f}, lift={sample['lift']:.4f}\")\n\n    return len(result.inserted_ids)\n\nn_rules_written = write_association_rules_to_mongo(spark_rules, db)","metadata":{"trusted":true,"execution":{"iopub.status.busy":"2026-04-06T16:55:37.327365Z","iopub.execute_input":"2026-04-06T16:55:37.327710Z","iopub.status.idle":"2026-04-06T16:55:38.507058Z","shell.execute_reply.started":"2026-04-06T16:55:37.327681Z","shell.execute_reply":"2026-04-06T16:55:38.506297Z"}},"outputs":[],"execution_count":null},{"cell_type":"code","source":"# ============================================================\n# CELL 25: WRITE USER RECOMMENDATIONS → MongoDB\n#\n# BIỆN LUẬN: Ở bước này dùng Popularity Baseline để tạo\n# top-12 cho 10K active users. Đây là \"cold start\" recommendation\n# — trước khi FP-Growth personalized model sẵn sàng, hệ thống\n# vẫn phải phục vụ được requests. Trong retail, không có\n# recommendation còn tệ hơn có recommendation phổ biến.\n# Time-decay popularity (recent weeks weighted higher) tốt hơn\n# pure popularity vì fashion trend thay đổi nhanh.\n# ============================================================\nimport numpy as np\n\ndef compute_time_decay_popularity(df_master, use_gpu=USE_GPU,\n                                   decay_weeks=4, top_n=50):\n    \"\"\"\n    Tính popularity score có time decay cho articles.\n    Score = Σ (count_in_week_w × exp(-λ × weeks_ago))\n    λ = log(2) / half_life_weeks\n\n    BIỆN LUẬN: Áo Hè ra mắt 2 tuần trước phải có score cao\n    hơn áo Đông của tháng trước dù áo Đông được mua nhiều hơn\n    tổng thể. Time decay giải quyết vấn đề này — exponential\n    decay là chuẩn mực trong industry (Netflix, Spotify dùng).\n    \"\"\"\n    print(\"  Computing time-decay popularity...\")\n\n    lambda_decay = math.log(2) / decay_weeks  # half-life = decay_weeks\n\n    if use_gpu:\n        df_pop = df_master[['article_id', 'week_of_year', 'days_since_start']].to_pandas()\n    else:\n        df_pop = df_master.select(['article_id', 'week_of_year',\n                                    'days_since_start']).to_pandas()\n\n    # Max week as reference\n    max_week = int(df_pop['week_of_year'].max())\n\n    # Group by article + week → count\n    weekly_counts = (\n        df_pop.groupby(['article_id', 'week_of_year'])\n        .size()\n        .reset_index(name='count')\n    )\n\n    # Tính weeks_ago\n    weekly_counts['weeks_ago'] = max_week - weekly_counts['week_of_year']\n    weekly_counts['weeks_ago'] = weekly_counts['weeks_ago'].clip(lower=0)\n\n    # Decay weight\n    weekly_counts['decay_weight'] = np.exp(-lambda_decay * weekly_counts['weeks_ago'])\n\n    # Weighted score\n    weekly_counts['weighted_count'] = (\n        weekly_counts['count'] * weekly_counts['decay_weight']\n    )\n\n    # Aggregate per article\n    pop_scores = (\n        weekly_counts.groupby('article_id')['weighted_count']\n        .sum()\n        .reset_index()\n        .rename(columns={'weighted_count': 'popularity_score'})\n        .sort_values('popularity_score', ascending=False)\n        .head(top_n)\n    )\n\n    pop_scores['rank'] = range(1, len(pop_scores) + 1)\n    pop_scores['article_id'] = pop_scores['article_id'].astype(str)\n\n    print(f\"  Top-{top_n} popular articles computed (time-decay λ={lambda_decay:.3f})\")\n    print(f\"  Top 5:\")\n    print(pop_scores.head(5).to_string(index=False))\n\n    return pop_scores\n\npop_df = compute_time_decay_popularity(df_master, decay_weeks=4, top_n=50)\n\n\ndef write_user_recommendations(db, top10k_path, pop_df, use_gpu=USE_GPU):\n    \"\"\"\n    Ghi top-12 recommendations cho 10K active users → MongoDB.\n    Dùng popularity baseline + user preference filter.\n    \"\"\"\n    print(\"\\n\" + \"=\" * 60)\n    print(\"SERVING: Write user recommendations → MongoDB\")\n    print(\"=\" * 60)\n\n    # Load top-10K active users (từ Chương 2)\n    if use_gpu:\n        import cudf\n        top10k = cudf.read_parquet(top10k_path).to_pandas()\n    else:\n        top10k = pl.read_parquet(top10k_path).to_pandas()\n\n    print(f\"  Users to process: {len(top10k):,}\")\n\n    # Load user features cho personalization\n    if use_gpu:\n        user_features = (\n            df_master[['customer_id', 'fav_colour', 'fav_department',\n                        'customer_price_tier', 'rfm_score', 'age_group',\n                        'fashion_season']]\n            .drop_duplicates('customer_id')\n            .to_pandas()\n        )\n    else:\n        user_features = (\n            df_master.select(['customer_id', 'fav_colour', 'fav_department',\n                               'customer_price_tier', 'rfm_score', 'age_group',\n                               'fashion_season'])\n            .unique('customer_id')\n            .to_pandas()\n        )\n\n    user_feat_dict = user_features.set_index('customer_id').to_dict(orient='index')\n\n    # Top-12 articles list từ pop_df (global baseline)\n    global_top12 = pop_df.head(12)\n\n    col    = db[COL_RECS]\n    ops    = []\n    now    = datetime.now(timezone.utc)\n    batch  = 500   # bulk write batch size\n\n    for i, row in top10k.iterrows():\n        cid  = str(row['customer_id'])\n        feat = user_feat_dict.get(cid, {})\n\n        # Build top-12 list (popularity baseline với metadata)\n        top12 = []\n        for _, art_row in global_top12.iterrows():\n            top12.append({\n                \"rank\"        : int(art_row['rank']),\n                \"article_id\"  : str(art_row['article_id']),\n                \"score\"       : round(float(art_row['popularity_score']), 4),\n                \"method\"      : \"time_decay_popularity\",\n                \"season_match\": feat.get('fashion_season', 'SS') == 'SS',\n            })\n\n        doc = {\n            \"customer_id\"    : cid,\n            \"rfm_score\"      : int(row.get('rfm_score', 0)),\n            \"customer_tier\"  : feat.get('customer_price_tier', 'mid'),\n            \"age_group\"      : feat.get('age_group', 'Unknown'),\n            \"fav_colour\"     : feat.get('fav_colour', 'Unknown'),\n            \"fav_department\" : feat.get('fav_department', 'Unknown'),\n            \"top12_articles\" : top12,\n            \"last_updated\"   : now,\n            \"week_generated\" : int(now.strftime('%W')),\n        }\n\n        # Upsert: update nếu tồn tại, insert nếu chưa có\n        ops.append(UpdateOne(\n            {\"customer_id\": cid},\n            {\"$set\": doc},\n            upsert=True\n        ))\n\n        # Flush batch\n        if len(ops) >= batch:\n            col.bulk_write(ops, ordered=False)\n            ops = []\n            print(f\"    Written {i+1:,} users...\", end='\\r')\n\n    # Flush remaining\n    if ops:\n        col.bulk_write(ops, ordered=False)\n\n    total = col.count_documents({})\n    print(f\"\\n  Total docs in user_recommendations: {total:,}\")\n\n    # Verify sample\n    sample = col.find_one({}, sort=[(\"rfm_score\", -1)])\n    print(f\"\\n  Sample user doc (highest RFM):\")\n    print(f\"    customer_id : {sample['customer_id']}\")\n    print(f\"    rfm_score   : {sample['rfm_score']}\")\n    print(f\"    tier        : {sample['customer_tier']}\")\n    print(f\"    top12 count : {len(sample['top12_articles'])}\")\n    print(f\"    top article : {sample['top12_articles'][0]}\")\n\ntop10k_path = f'{OUTPUT_DIR}/top10k_active_users.parquet'\nwrite_user_recommendations(db, top10k_path, pop_df)","metadata":{"trusted":true,"execution":{"iopub.status.busy":"2026-04-06T16:55:48.397760Z","iopub.execute_input":"2026-04-06T16:55:48.398323Z","iopub.status.idle":"2026-04-06T16:57:27.266933Z","shell.execute_reply.started":"2026-04-06T16:55:48.398291Z","shell.execute_reply":"2026-04-06T16:57:27.266258Z"}},"outputs":[],"execution_count":null},{"cell_type":"code","source":"# ============================================================\n# CELL 26: CRUD DEMO — MongoDB tương tác đầy đủ\n#\n# ĐÂY LÀ PHẦN QUAN TRỌNG ĐỂ MINH CHỨNG VỚI GIẢNG VIÊN\n# Thực hiện 5 operations: Create, Read, Update, Delete + Agg\n# ============================================================\n\nprint(\"=\" * 60)\nprint(\"CRUD DEMO — MongoDB Atlas\")\nprint(\"=\" * 60)\n\ncol_recs  = db[COL_RECS]\ncol_rules = db[COL_RULES]\n\n# ─────────────────────────────────────────────\n# C — CREATE: Thêm user mới (new signup)\n# ─────────────────────────────────────────────\nprint(\"\\n[C] CREATE: Thêm user mới vừa đăng ký\")\n\nnew_user = {\n    \"customer_id\"    : \"DEMO_USER_001\",\n    \"rfm_score\"      : 3,        # new user → lowest RFM\n    \"customer_tier\"  : \"mid\",\n    \"age_group\"      : \"Millennial\",\n    \"fav_colour\"     : \"Black\",  # default preference\n    \"fav_department\" : \"Unknown\",\n    \"top12_articles\" : [         # cold-start: dùng global popular items\n        {\"rank\": i+1, \"article_id\": str(pop_df.iloc[i]['article_id']),\n         \"score\": round(float(pop_df.iloc[i]['popularity_score']), 4),\n         \"method\": \"cold_start_popular\", \"season_match\": True}\n        for i in range(12)\n    ],\n    \"last_updated\"   : datetime.now(timezone.utc),\n    \"week_generated\" : int(datetime.now().strftime('%W')),\n    \"is_demo\"        : True,\n}\n\nresult = col_recs.update_one(\n    {\"customer_id\": \"DEMO_USER_001\"},\n    {\"$set\": new_user},\n    upsert=True\n)\nprint(f\"     Upserted: matched={result.matched_count}, modified={result.modified_count}\")\nprint(f\"     upserted_id = {result.upserted_id}\")\n\n\n# ─────────────────────────────────────────────\n# R — READ: Truy vấn recommendation\n# ─────────────────────────────────────────────\nprint(\"\\n[R] READ: Lấy recommendation của DEMO_USER_001\")\n\ndoc = col_recs.find_one(\n    {\"customer_id\": \"DEMO_USER_001\"},\n    {\"_id\": 0, \"customer_id\": 1, \"rfm_score\": 1,\n     \"customer_tier\": 1, \"top12_articles\": {\"$slice\": 3}}\n)\nprint(f\"  customer_id : {doc['customer_id']}\")\nprint(f\"  rfm_score   : {doc['rfm_score']}\")\nprint(f\"  tier        : {doc['customer_tier']}\")\nprint(f\"  top 3 recs  :\")\nfor art in doc['top12_articles']:\n    print(f\"    rank={art['rank']} article={art['article_id']} score={art['score']}\")\n\n\n# ─────────────────────────────────────────────\n# U — UPDATE: User mua thêm → cập nhật profile\n# ─────────────────────────────────────────────\nprint(\"\\n[U] UPDATE: Cập nhật sau khi user mua thêm\")\n\n# Giả lập: user vừa mua 5 items → rfm_score tăng, tier lên premium\nnew_purchases = [\"0572024001\", \"0610776002\", \"0685813001\"]\n\nupdate_result = col_recs.update_one(\n    {\"customer_id\": \"DEMO_USER_001\"},\n    {\n        \"$set\" : {\n            \"rfm_score\"     : 8,\n            \"customer_tier\" : \"premium\",\n            \"fav_colour\"    : \"Dark Blue\",\n            \"last_updated\"  : datetime.now(timezone.utc),\n        },\n        \"$push\": {\n            \"purchase_history\": {\n                \"$each\"     : new_purchases,\n                \"$slice\"    : -50,      # giữ tối đa 50 items gần nhất\n                \"$position\" : 0,        # thêm vào đầu\n            }\n        },\n        \"$inc\" : {\"total_purchases\": len(new_purchases)},\n    }\n)\nprint(f\"  Updated: matched={update_result.matched_count}, modified={update_result.modified_count}\")\n\n# Verify\nupdated = col_recs.find_one(\n    {\"customer_id\": \"DEMO_USER_001\"},\n    {\"_id\": 0, \"rfm_score\": 1, \"customer_tier\": 1,\n     \"fav_colour\": 1, \"total_purchases\": 1}\n)\nprint(f\"  New rfm_score   : {updated['rfm_score']}\")\nprint(f\"  New tier        : {updated['customer_tier']}\")\nprint(f\"  New fav_colour  : {updated['fav_colour']}\")\nprint(f\"  total_purchases : {updated.get('total_purchases', 'N/A')}\")\n\n\n# ─────────────────────────────────────────────\n# D — DELETE: Xóa inactive user (GDPR compliance)\n# ─────────────────────────────────────────────\nprint(\"\\n[D] DELETE: Xóa DEMO_USER_001 (GDPR request)\")\n\n# Trước khi xóa — audit log\naudit_entry = {\n    \"action\"      : \"gdpr_delete\",\n    \"customer_id\" : \"DEMO_USER_001\",\n    \"deleted_at\"  : datetime.now(timezone.utc),\n    \"reason\"      : \"user_request\",\n}\ndb[COL_META].insert_one(audit_entry)\n\n# Xóa user\ndel_result = col_recs.delete_one({\"customer_id\": \"DEMO_USER_001\"})\nprint(f\"     Deleted: {del_result.deleted_count} document(s)\")\nprint(f\"     Audit log written to {COL_META}\")\n\n# Verify xóa thành công\nverify = col_recs.find_one({\"customer_id\": \"DEMO_USER_001\"})\nprint(f\"  Verify: {'NOT FOUND' if verify is None else 'STILL EXISTS'}\")\n\n\n# ─────────────────────────────────────────────\n# A — AGGREGATE: Business insight query\n# ─────────────────────────────────────────────\nprint(\"\\n[A] AGGREGATE: Phân phối tier của 10K users\")\n\npipeline = [\n    {\"$group\": {\n        \"_id\"   : \"$customer_tier\",\n        \"count\" : {\"$sum\": 1},\n        \"avg_rfm\": {\"$avg\": \"$rfm_score\"},\n    }},\n    {\"$sort\": {\"count\": -1}},\n    {\"$project\": {\n        \"_id\"    : 0,\n        \"tier\"   : \"$_id\",\n        \"count\"  : 1,\n        \"avg_rfm\": {\"$round\": [\"$avg_rfm\", 2]},\n    }}\n]\n\ntier_dist = list(col_recs.aggregate(pipeline))\nprint(f\"  {'Tier':<12} {'Count':>8} {'Avg RFM':>10}\")\nprint(f\"  {'-'*32}\")\nfor row in tier_dist:\n    print(f\"  {row['tier']:<12} {row['count']:>8,} {row['avg_rfm']:>10.2f}\")\n\n\n# ─────────────────────────────────────────────\n# Kiểm tra dung lượng MongoDB\n# ─────────────────────────────────────────────\nprint(\"\\n[INFO] MongoDB storage check\")\nstats = db.command(\"dbStats\")\nsize_mb = stats['storageSize'] / 1e6\ndata_mb = stats['dataSize'] / 1e6\nprint(f\"  Data size    : {data_mb:.2f} MB\")\nprint(f\"  Storage size : {size_mb:.2f} MB\")\nprint(f\"  Budget used  : {size_mb/500*100:.1f}% of 500MB free tier\")\nprint(f\"  Collections  :\")\nfor col_name in db.list_collection_names():\n    col_stat = db.command(\"collStats\", col_name)\n    print(f\"    {col_name:<30}: {col_stat['storageSize']/1e6:.3f} MB \"\n          f\"({col_stat['count']:,} docs)\")","metadata":{"trusted":true,"execution":{"iopub.status.busy":"2026-04-06T16:57:49.953269Z","iopub.execute_input":"2026-04-06T16:57:49.953537Z","iopub.status.idle":"2026-04-06T16:57:53.541345Z","shell.execute_reply.started":"2026-04-06T16:57:49.953514Z","shell.execute_reply":"2026-04-06T16:57:53.540615Z"}},"outputs":[],"execution_count":null},{"cell_type":"code","source":"# ============================================================\n# CELL 27: SPEED LAYER — Real-time item trend update\n#\n# BIỆN LUẬN: Speed layer xử lý events real-time (từ Kafka)\n# để cập nhật popularity nhanh hơn batch cycle (daily).\n# Khi H&M launch collection mới, items cần xuất hiện trong\n# recommendation trong vài phút, không phải ngày hôm sau.\n# Ta mock bằng việc process kafka_events đã thu thập ở Cell 22.\n# ============================================================\nfrom collections import defaultdict\n\ndef process_speed_layer(kafka_events, db, top_n=20):\n    \"\"\"\n    Xử lý events từ Kafka, cập nhật item_trends collection\n    trong MongoDB theo real-time (giả lập).\n    \"\"\"\n    print(\"=\" * 60)\n    print(\"SPEED LAYER: Real-time trend update\")\n    print(\"=\" * 60)\n\n    # Đếm purchase frequency từ Kafka events\n    article_counts = defaultdict(int)\n    channel_counts = defaultdict(lambda: defaultdict(int))\n    season_counts  = defaultdict(lambda: defaultdict(int))\n\n    for event in kafka_events:\n        aid     = str(event.get('article_id', ''))\n        channel = event.get('channel', 1)\n        season  = event.get('season', 'SS')\n        article_counts[aid] += 1\n        channel_counts[aid][f'ch_{channel}'] += 1\n        season_counts[aid][season] += 1\n\n    # Top N trending articles\n    top_trending = sorted(\n        article_counts.items(),\n        key=lambda x: x[1],\n        reverse=True\n    )[:top_n]\n\n    print(f\"  Events processed  : {len(kafka_events):,}\")\n    print(f\"  Unique articles   : {len(article_counts):,}\")\n    print(f\"  Top 5 trending    :\")\n    for aid, cnt in top_trending[:5]:\n        print(f\"    article {aid}: {cnt:,} purchases\")\n\n    # Upsert vào MongoDB item_trends\n    col   = db[COL_TRENDS]\n    now   = datetime.now(timezone.utc)\n    week  = int(now.strftime('%W'))\n    ops   = []\n\n    for aid, cnt in top_trending:\n        doc = {\n            \"article_id\"      : aid,\n            \"week\"            : week,\n            \"purchase_count\"  : cnt,\n            \"channel_dist\"    : dict(channel_counts[aid]),\n            \"season_dist\"     : dict(season_counts[aid]),\n            \"trending_score\"  : cnt,        # simplification: real = decay score\n            \"updated_at\"      : now,\n        }\n        ops.append(UpdateOne(\n            {\"article_id\": aid, \"week\": week},\n            {\"$set\": doc},\n            upsert=True\n        ))\n\n    if ops:\n        result = col.bulk_write(ops, ordered=False)\n        print(f\"\\n MongoDB item_trends: {result.upserted_count} inserted, \"\n              f\"{result.modified_count} updated\")\n\n    # Query trending để verify\n    trending_sample = list(\n        col.find({}, {\"_id\": 0, \"article_id\": 1, \"purchase_count\": 1,\n                      \"season_dist\": 1})\n           .sort(\"trending_score\", -1)\n           .limit(5)\n    )\n    print(f\"\\n  Current trending (top 5):\")\n    for item in trending_sample:\n        print(f\"    {item['article_id']}: {item['purchase_count']} purchases | \"\n              f\"season={item['season_dist']}\")\n\nprocess_speed_layer(kafka_events, db, top_n=20)","metadata":{"trusted":true,"execution":{"iopub.status.busy":"2026-04-06T16:57:58.630500Z","iopub.execute_input":"2026-04-06T16:57:58.631083Z","iopub.status.idle":"2026-04-06T16:57:59.207783Z","shell.execute_reply.started":"2026-04-06T16:57:58.631055Z","shell.execute_reply":"2026-04-06T16:57:59.207014Z"}},"outputs":[],"execution_count":null},{"cell_type":"markdown","source":"# **CHƯƠNG 3: Máy học**","metadata":{}},{"cell_type":"markdown","source":"## **Giai Đoạn 3.1: Time-based Split + FP-Growth Training**","metadata":{}},{"cell_type":"code","source":"# ============================================================\n# CELL 28: TIME-BASED TRAIN / TEST SPLIT\n#\n# BIỆN LUẬN CHỐNG DATA LEAKAGE:\n# Random split là sai hoàn toàn trong bài toán recommendation\n# có yếu tố thời gian vì:\n#   (1) Transaction ngày thứ Sáu có thể rơi vào train,\n#       transaction ngày thứ Hai cùng tuần rơi vào test →\n#       model \"nhìn thấy tương lai\" khi học patterns.\n#   (2) Association rules được khai phá từ baskets đã bao gồm\n#       cả items của test period → support/confidence bị inflate.\n#   (3) MAP@12 sẽ cao giả tạo: predict items mà model đã thấy\n#       trong \"future\" transactions.\n#\n# Time-based split: train = tất cả dữ liệu trước cutoff,\n# test = 7 ngày cuối cùng. Đây là cách H&M Kaggle competition\n# đánh giá chính thức (predict last week).\n# ============================================================\nimport pandas as pd\nimport numpy as np\nimport gc\n\ndef time_based_split(df_master, test_days=7, use_gpu=USE_GPU):\n    \"\"\"\n    Chia dữ liệu theo thời gian.\n    Train: [start, cutoff)\n    Test:  [cutoff, end]\n    \"\"\"\n    print(\"=\" * 60)\n    print(\"TIME-BASED TRAIN/TEST SPLIT\")\n    print(\"=\" * 60)\n\n    # Lấy ngày tmax\n    if use_gpu:\n        import cudf\n        dates = cudf.to_datetime(df_master['t_dat'])\n        t_max = dates.max()\n        t_min = dates.min()\n        \n        # SỬA LỖI Ở ĐÂY: Dùng pd.Timedelta cho tính toán giá trị đơn (scalar)\n        cutoff = t_max - pd.Timedelta(days=test_days)\n\n        df_master['_date'] = dates\n        train_mask = df_master['_date'] < cutoff\n        test_mask  = df_master['_date'] >= cutoff\n\n        df_train = df_master[train_mask].drop('_date', axis=1)\n        df_test  = df_master[test_mask].drop('_date', axis=1)\n\n        t_min_str    = str(t_min)[:10]\n        cutoff_str   = str(cutoff)[:10]\n        t_max_str    = str(t_max)[:10]\n\n    else:\n        import polars as pl\n\n        df_master = df_master.with_columns(\n            pl.col('t_dat').str.to_date().alias('_date')\n        )\n        t_max  = df_master['_date'].max()\n        t_min  = df_master['_date'].min()\n        import datetime\n        cutoff = t_max - datetime.timedelta(days=test_days)\n\n        df_train = (\n            df_master.filter(pl.col('_date') < cutoff)\n            .drop('_date')\n        )\n        df_test = (\n            df_master.filter(pl.col('_date') >= cutoff)\n            .drop('_date')\n        )\n\n        t_min_str  = str(t_min)\n        cutoff_str = str(cutoff)\n        t_max_str  = str(t_max)\n\n    print(f\"  Dataset range : {t_min_str} → {t_max_str}\")\n    print(f\"  Cutoff        : {cutoff_str} (last {test_days} days = TEST)\")\n    print(f\"\\n  TRAIN  : {len(df_train):>12,} rows\")\n    print(f\"  TEST   : {len(df_test):>12,} rows\")\n    print(f\"  Ratio  : {len(df_train)/len(df_master)*100:.1f}% / \"\n          f\"{len(df_test)/len(df_master)*100:.1f}%\")\n\n    # Test set: unique customers\n    if use_gpu:\n        n_test_cust  = df_test['customer_id'].nunique()\n        n_train_cust = df_train['customer_id'].nunique()\n    else:\n        n_test_cust  = df_test['customer_id'].n_unique()\n        n_train_cust = df_train['customer_id'].n_unique()\n\n    print(f\"\\n  Train customers : {n_train_cust:,}\")\n    print(f\"  Test customers  : {n_test_cust:,}\")\n\n    # Ground truth: {customer_id → set of articles bought in test week}\n    print(f\"\\n  Building ground truth dict...\")\n    if use_gpu:\n        gt_pd = df_test[['customer_id', 'article_id']].to_pandas()\n    else:\n        gt_pd = df_test.select(['customer_id', 'article_id']).to_pandas()\n\n    gt_pd['article_id'] = gt_pd['article_id'].astype(str)\n    ground_truth = (\n        gt_pd.groupby('customer_id')['article_id']\n        .apply(set)\n        .to_dict()\n    )\n    print(f\"  Ground truth entries: {len(ground_truth):,}\")\n\n    # Avg items per customer in test\n    avg_gt = np.mean([len(v) for v in ground_truth.values()])\n    print(f\"  Avg items/customer in test: {avg_gt:.1f}\")\n\n    return df_train, df_test, ground_truth, cutoff_str\n\ndf_train, df_test, ground_truth, cutoff_str = time_based_split(\n    df_master, test_days=7\n)\ngc.collect()\nprint(\"\\nSplit complete — no data leakage\")","metadata":{"trusted":true,"execution":{"iopub.status.busy":"2026-04-06T17:03:49.518858Z","iopub.execute_input":"2026-04-06T17:03:49.519417Z","iopub.status.idle":"2026-04-06T17:04:18.426155Z","shell.execute_reply.started":"2026-04-06T17:03:49.519387Z","shell.execute_reply":"2026-04-06T17:04:18.425619Z"}},"outputs":[],"execution_count":null},{"cell_type":"code","source":"# ============================================================\n# CELL 29: BUILD TRAINING BASKETS\n#\n# BIỆN LUẬN: Basket = tất cả items một customer mua trong\n# toàn bộ train period (không phải per-transaction).\n# Lý do gộp theo customer thay vì per-session:\n#   - H&M không có session ID trong dataset\n#   - \"Customer lifetime basket\" capture long-term preference\n#     tốt hơn cho fashion (người ta mua outfit theo style,\n#     không chỉ theo 1 session)\n#   - Cách này align với H&M Kaggle top solutions\n# ============================================================\n\ndef build_training_baskets(df_train, use_gpu=USE_GPU,\n                            min_basket_size=2,\n                            max_basket_size=100):\n    \"\"\"\n    Tạo customer baskets từ train set.\n    Lọc: buyer_type='normal', basket size [2, 100].\n    \"\"\"\n    print(\"=\" * 60)\n    print(\"BUILD TRAINING BASKETS\")\n    print(\"=\" * 60)\n\n    # Lọc normal buyers (kèm cơ chế an toàn nếu thiếu cột)\n    if use_gpu:\n        if 'buyer_type' in df_train.columns:\n            df_normal = df_train[df_train['buyer_type'] == 'normal']\n        else:\n            df_normal = df_train\n        df_pd = df_normal[['customer_id', 'article_id']].to_pandas()\n    else:\n        import polars as pl\n        if 'buyer_type' in df_train.columns:\n            df_pd = (\n                df_train.filter(pl.col('buyer_type') == 'normal')\n                .select(['customer_id', 'article_id'])\n                .to_pandas()\n            )\n        else:\n            df_pd = df_train.select(['customer_id', 'article_id']).to_pandas()\n\n    df_pd['article_id'] = df_pd['article_id'].astype(str)\n\n    # Build baskets\n    baskets_raw = (\n        df_pd.groupby('customer_id')['article_id']\n        .apply(lambda x: list(x.unique()))\n    )\n\n    # Filter by size\n    baskets = baskets_raw[\n        (baskets_raw.apply(len) >= min_basket_size) &\n        (baskets_raw.apply(len) <= max_basket_size)\n    ].reset_index()\n    baskets.columns = ['customer_id', 'items']\n\n    sizes = baskets['items'].apply(len)\n    print(f\"  Raw customers    : {len(baskets_raw):,}\")\n    print(f\"  After size filter: {len(baskets):,}\")\n    print(f\"  Basket size stats:\")\n    print(f\"    min = {sizes.min()}\")\n    print(f\"    max = {sizes.max()}\")\n    print(f\"    mean= {sizes.mean():.1f}\")\n    print(f\"    p50 = {sizes.quantile(0.5):.0f}\")\n    print(f\"    p90 = {sizes.quantile(0.9):.0f}\")\n    print(f\"    p99 = {sizes.quantile(0.99):.0f}\")\n\n    # Sample basket\n    print(f\"\\n  Sample basket (3 items shown):\")\n    sample = baskets.iloc[0]\n    print(f\"    customer: {sample['customer_id']}\")\n    print(f\"    items   : {sample['items'][:3]}...\")\n\n    return baskets\n\ntrain_baskets = build_training_baskets(df_train)\nprint(f\"\\nTraining baskets ready: {len(train_baskets):,}\")","metadata":{"trusted":true,"execution":{"iopub.status.busy":"2026-04-06T17:05:22.436465Z","iopub.execute_input":"2026-04-06T17:05:22.437055Z","iopub.status.idle":"2026-04-06T17:06:34.102962Z","shell.execute_reply.started":"2026-04-06T17:05:22.437024Z","shell.execute_reply":"2026-04-06T17:06:34.102202Z"}},"outputs":[],"execution_count":null},{"cell_type":"code","source":"# ============================================================\n# CELL 30: FP-GROWTH TRAINING — mlxtend + Hyperparameter sweep\n# BIỆN LUẬN dùng mlxtend trên CPU thay vì Spark MLlib:\n# - mlxtend FP-Growth tối ưu cho single-machine, dễ tune\n# - Spark MLlib hợp hơn khi cần scale horizontal (multi-node)\n# - Trên Kaggle T4, bottleneck là memory, không phải compute\n# - mlxtend cho phép truy cập trực tiếp frequent itemset\n#   → dễ tính Conviction, Leverage ngoài Support/Confidence/Lif\n# Grid search ngưỡng để tìm trade-off tối ưu:\n# - min_support quá thấp → quá nhiều rules rác, chậm\n# - min_support quá cao → bỏ sót cross-category insights\n# ============================================================\nimport warnings\n# Bỏ qua mọi cảnh báo khó chịu từ các thư viện cũ\nwarnings.filterwarnings('ignore', category=DeprecationWarning)\n\nfrom mlxtend.frequent_patterns import fpgrowth, association_rules\nfrom mlxtend.preprocessing import TransactionEncoder\nimport pandas as pd\nimport numpy as np\nimport time\n\ndef encode_baskets(baskets_df, sample_n=None):\n    \"\"\"\n    Chuyển list-of-items → one-hot encoded DataFrame.\n    BIỆN LUẬN: TransactionEncoder tạo sparse boolean matrix —\n    hiệu quả hơn dense float matrix vì basket data rất sparse\n    (khách thường mua <1% catalog articles).\n    \"\"\"\n    print(\"  Encoding baskets → one-hot matrix...\")\n    t0 = time.time()\n\n    items_list = baskets_df['items'].tolist()\n    if sample_n and sample_n < len(items_list):\n        import random\n        random.seed(42)\n        items_list = random.sample(items_list, sample_n)\n        print(f\"  Sampled {sample_n:,} baskets từ {len(baskets_df):,}\")\n\n    te = TransactionEncoder()\n    te_array = te.fit(items_list).transform(items_list, sparse=True)\n\n    # Sparse → DataFrame (sparse columns để tiết kiệm RAM)\n    df_encoded = pd.DataFrame.sparse.from_spmatrix(\n        te_array,\n        columns=te.columns_\n    )\n\n    elapsed = time.time() - t0\n    density = te_array.nnz / (te_array.shape[0] * te_array.shape[1])\n    print(f\"  Shape    : {df_encoded.shape}\")\n    print(f\"  Density  : {density:.4f} ({density*100:.2f}%) — very sparse ✓\")\n    print(f\"  Elapsed  : {elapsed:.2f}s\")\n\n    return df_encoded, te.columns_\n\n# Sample 80K baskets để FP-Growth chạy nhanh trên Kaggle\n# Đây vẫn đủ statistical power (~6% total train customers)\ndf_encoded, article_cols = encode_baskets(train_baskets, sample_n=80_000)\n\n\ndef run_fpgrowth_sweep(df_encoded, support_values, confidence_values):\n    \"\"\"\n    Grid search trên min_support × min_confidence.\n    Tìm cấu hình cho số rules nhiều nhất với lift > 1.\n    \"\"\"\n    print(\"\\n\" + \"=\" * 60)\n    print(\"FP-GROWTH HYPERPARAMETER SWEEP\")\n    print(\"=\" * 60)\n\n    results = []\n\n    for sup in support_values:\n        t0 = time.time()\n        # Bước 1: Frequent itemsets\n        freq_items = fpgrowth(\n            df_encoded,\n            min_support=sup,\n            use_colnames=True,\n            max_len=3,       # tối đa 3-itemset (đủ cho fashion rules)\n            verbose=0,\n        )\n        elapsed_freq = time.time() - t0\n\n        n_freq = len(freq_items)\n\n        for conf in confidence_values:\n            if n_freq == 0:\n                results.append({\n                    'min_support': sup, 'min_confidence': conf,\n                    'n_freq_itemsets': 0, 'n_rules': 0,\n                    'n_rules_lift_gt1': 0, 'elapsed_s': elapsed_freq,\n                })\n                continue\n\n            t1 = time.time()\n            rules = association_rules(\n                freq_items,\n                metric=\"confidence\",\n                min_threshold=conf,\n                num_itemsets=n_freq,\n            )\n            elapsed_rules = time.time() - t1\n\n            # Thêm Conviction\n            rules['conviction'] = (\n                (1 - rules['consequent support']) /\n                (1 - rules['confidence'] + 1e-9)\n            ).round(4)\n\n            n_lift = int((rules['lift'] > 1.0).sum())\n\n            results.append({\n                'min_support'      : sup,\n                'min_confidence'   : conf,\n                'n_freq_itemsets'  : n_freq,\n                'n_rules'          : len(rules),\n                'n_rules_lift_gt1' : n_lift,\n                'elapsed_s'        : round(elapsed_freq + elapsed_rules, 2),\n            })\n\n            print(f\"  sup={sup:.3f} conf={conf:.2f} → \"\n                  f\"freq={n_freq:>5,} rules={len(rules):>5,} \"\n                  f\"lift>1={n_lift:>5,}  ({elapsed_freq+elapsed_rules:.1f}s)\")\n\n    results_df = pd.DataFrame(results)\n    return results_df\n\n# Grid: 3 support × 2 confidence = 6 runs\n# FIX LỖI 0 RULES: Hạ ngưỡng support xuống thấp như bản PySpark\nsweep_results = run_fpgrowth_sweep(\n    df_encoded,\n    support_values    = [0.002, 0.003, 0.005], # 0.002 trên 20K giỏ ~ 40 lần xuất hiện\n    confidence_values = [0.05, 0.1],           \n)\n\nprint(\"\\n  Sweep results:\")\nprint(sweep_results.to_string(index=False))\n\n# Chọn cấu hình tốt nhất: max rules với lift > 1\nbest_idx = sweep_results['n_rules_lift_gt1'].idxmax()\nbest_cfg = sweep_results.iloc[best_idx]\nprint(f\"\\n  Best config:\")\nprint(f\"    min_support    = {best_cfg['min_support']}\")\nprint(f\"    min_confidence = {best_cfg['min_confidence']}\")\nprint(f\"    n_rules(lift>1)= {best_cfg['n_rules_lift_gt1']:.0f}\")","metadata":{"trusted":true,"execution":{"iopub.status.busy":"2026-04-06T17:08:19.150878Z","iopub.execute_input":"2026-04-06T17:08:19.151166Z","execution_failed":"2026-04-06T17:58:46.785Z"}},"outputs":[],"execution_count":null},{"cell_type":"code","source":"# ============================================================\n# CELL 31: TRAIN FINAL FP-GROWTH MODEL\n# ============================================================\n\ndef train_final_fpgrowth(df_encoded, min_support, min_confidence):\n    \"\"\"\n    Train FP-Growth với cấu hình đã chọn.\n    Tính đầy đủ: Support, Confidence, Lift, Conviction, Leverage.\n    \"\"\"\n    print(\"=\" * 60)\n    print(\"FINAL FP-GROWTH TRAINING\")\n    print(\"=\" * 60)\n    print(f\"  min_support    = {min_support}\")\n    print(f\"  min_confidence = {min_confidence}\")\n\n    t0 = time.time()\n\n    # Frequent itemsets\n    freq_items = fpgrowth(\n        df_encoded,\n        min_support=min_support,\n        use_colnames=True,\n        max_len=3,\n        verbose=1,\n    )\n\n    # Association rules\n    rules = association_rules(\n        freq_items,\n        metric=\"confidence\",\n        min_threshold=min_confidence,\n        num_itemsets=len(freq_items),\n    )\n\n    elapsed = time.time() - t0\n    print(f\"\\n  Training time   : {elapsed:.1f}s\")\n    print(f\"  Frequent sets   : {len(freq_items):,}\")\n    print(f\"  Total rules     : {len(rules):,}\")\n\n    # ── Tính thêm metrics ──\n\n    # Conviction: P(not C) / P(not C | A)\n    # Cao → rule rất đáng tin (A luôn kéo theo C)\n    rules['conviction'] = (\n        (1 - rules['consequent support']) /\n        (1 - rules['confidence'] + 1e-9)\n    ).round(6)\n\n    # Leverage: P(A∩C) - P(A)×P(C)\n    # Dương → A và C xảy ra cùng nhau nhiều hơn expected\n    rules['leverage'] = (\n        rules['support'] -\n        rules['antecedent support'] * rules['consequent support']\n    ).round(6)\n\n    # Lọc lift > 1 (positive association)\n    rules_positive = rules[rules['lift'] > 1.0].copy()\n    print(f\"  Rules (lift>1)  : {len(rules_positive):,}\")\n\n    # Convert frozenset → list/string cho readable output\n    rules_positive['antecedents_str'] = rules_positive['antecedents'].apply(\n        lambda x: list(x)\n    )\n    rules_positive['consequents_str'] = rules_positive['consequents'].apply(\n        lambda x: list(x)\n    )\n\n    # Sort by lift descending\n    rules_positive = rules_positive.sort_values('lift', ascending=False)\n\n    # ── Metrics Summary ──\n    print(f\"\\n  Metrics summary (rules với lift>1):\")\n    metrics = ['support', 'confidence', 'lift', 'conviction', 'leverage']\n    for m in metrics:\n        col_data = rules_positive[m]\n        print(f\"    {m:12s}: mean={col_data.mean():.4f}  \"\n              f\"min={col_data.min():.4f}  max={col_data.max():.4f}\")\n\n    # Top 20 rules\n    print(f\"\\n  Top 10 rules by lift:\")\n    display_cols = ['antecedents_str', 'consequents_str',\n                    'support', 'confidence', 'lift', 'conviction']\n    print(rules_positive[display_cols].head(10).to_string(index=False))\n\n    return rules_positive, freq_items\n\n# Dùng best config từ sweep\nBEST_SUPPORT    = float(best_cfg['min_support'])\nBEST_CONFIDENCE = float(best_cfg['min_confidence'])\n\nfp_rules, fp_freq = train_final_fpgrowth(\n    df_encoded,\n    min_support    = BEST_SUPPORT,\n    min_confidence = BEST_CONFIDENCE,\n)","metadata":{"trusted":true},"outputs":[],"execution_count":null},{"cell_type":"code","source":"# ============================================================\n# CELL 32: CANDIDATE GENERATION\n#\n# BIỆN LUẬN — Tại sao cần candidate generation?\n# FP-Growth tạo ra rules dạng {A} → {B}. Để tạo top-12\n# recommendation cho user, ta không thể chỉ lookup rules —\n# nhiều users chưa mua item nào có rules.\n# Kaggle Top Solutions dùng multi-source candidate pool:\n#   (1) FP-Growth rules (personalized)\n#   (2) Time-decay popularity (cold start + fallback)\n#   (3) Repurchase candidates (items user đã mua trước đây)\n# Merge 3 sources → rank → lấy top-12.\n# Đây là \"cascade recommendation\" — precision cao hơn\n# single-source vì mỗi source cover weakness của source khác.\n# ============================================================\n\ndef generate_candidates(df_train, fp_rules, pop_df,\n                         customer_ids=None, top_k=50,\n                         use_gpu=USE_GPU):\n    \"\"\"\n    Tạo candidate pool cho mỗi customer từ 3 sources.\n    Trả về dict: {customer_id → [(article_id, score, source), ...]}\n    \"\"\"\n    print(\"=\" * 60)\n    print(\"CANDIDATE GENERATION (3 sources)\")\n    print(\"=\" * 60)\n\n    # ── Source 1: FP-Growth rule lookup ──\n    # {antecedent_item → [(consequent, confidence, lift), ...]}\n    rule_lookup = {}\n    for _, row in fp_rules.iterrows():\n        for ant in row['antecedents']:\n            if ant not in rule_lookup:\n                rule_lookup[ant] = []\n            for cons in row['consequents']:\n                rule_lookup[ant].append((\n                    cons,\n                    float(row['confidence']),\n                    float(row['lift']),\n                ))\n    print(f\"  Source 1 (FP-Growth): {len(rule_lookup):,} antecedent items mapped\")\n\n    # ── Source 2: Global popularity (time-decay) ──\n    pop_articles = pop_df['article_id'].astype(str).tolist()\n    pop_scores   = pop_df['popularity_score'].tolist()\n    pop_max      = max(pop_scores) if pop_scores else 1\n    pop_norm     = {a: s/pop_max for a, s in zip(pop_articles, pop_scores)}\n    print(f\"  Source 2 (Popularity): top {len(pop_articles)} articles\")\n\n    # ── Source 3: Repurchase history per customer ──\n    if use_gpu:\n        repurchase_pd = df_train[['customer_id', 'article_id']].to_pandas()\n    else:\n        import polars as pl\n        repurchase_pd = df_train.select(['customer_id', 'article_id']).to_pandas()\n\n    repurchase_pd['article_id'] = repurchase_pd['article_id'].astype(str)\n    repurchase_hist = (\n        repurchase_pd.groupby('customer_id')['article_id']\n        .apply(set)\n        .to_dict()\n    )\n    print(f\"  Source 3 (Repurchase): {len(repurchase_hist):,} customers\")\n\n    # ── Generate candidates per customer ──\n    if customer_ids is None:\n        # Dùng test set customers\n        customer_ids = list(ground_truth.keys())[:5_000]  # 5K để nhanh\n\n    print(f\"\\n  Generating for {len(customer_ids):,} customers...\")\n\n    candidates = {}\n    for cid in customer_ids:\n        scored = {}\n\n        # Source 3: repurchase (highest prior)\n        history = repurchase_hist.get(cid, set())\n        for art in list(history)[:20]:   # top 20 recent\n            scored[art] = scored.get(art, 0) + 0.6   # repurchase weight\n\n        # Source 1: FP-Growth rules from history\n        for art in list(history)[:10]:   # apply rules on top 10 items\n            if art in rule_lookup:\n                for cons, conf, lift in rule_lookup[art][:5]:\n                    fp_score = conf * min(lift / 5, 1.0)  # normalize lift\n                    scored[cons] = scored.get(cons, 0) + fp_score * 0.8\n\n        # Source 2: popularity fallback (for all users)\n        for art, pop_score in zip(pop_articles[:top_k], [pop_norm[a] for a in pop_articles[:top_k]]):\n            if art not in scored:\n                scored[art] = pop_score * 0.4  # lower weight for generic\n\n        # Exclude items already in history (avoid trivial repurchase in eval)\n        for art in history:\n            scored.pop(art, None)\n\n        # Rank by score, take top_k\n        top_candidates = sorted(scored.items(), key=lambda x: x[1], reverse=True)[:top_k]\n        candidates[cid] = top_candidates\n\n    print(f\"  Candidates generated\")\n    avg_cands = np.mean([len(v) for v in candidates.values()])\n    print(f\"  Avg candidates/user: {avg_cands:.1f}\")\n\n    return candidates\n\ncandidates = generate_candidates(\n    df_train, fp_rules, pop_df,\n    customer_ids=list(ground_truth.keys())[:5_000],\n    top_k=50,\n)","metadata":{"trusted":true},"outputs":[],"execution_count":null},{"cell_type":"code","source":"# ============================================================\n# CELL 33: MAP@12 — METRIC ĐÁNH GIÁ\n#\n# BIỆN LUẬN — Tại sao MAP@12?\n# H&M Kaggle competition dùng MAP@12 vì:\n#   (1) Customers thường không scroll quá 12 items trên mobile\n#   (2) MAP (Mean Average Precision) phạt nặng nếu item đúng\n#       xếp hạng thấp → khuyến khích đưa item đúng lên đầu\n#   (3) @12 là industry standard cho \"above the fold\" trên app\n#\n# MAP@K = (1/|U|) × Σ_u AP@K(u)\n# AP@K(u) = (1/min(|Rel_u|, K)) × Σ_{k=1}^{K} P@k × rel(k)\n# P@k = precision tại vị trí k\n# rel(k) = 1 nếu item tại vị trí k là relevant, else 0\n# ============================================================\n\ndef average_precision_at_k(recommended, relevant, k=12):\n    \"\"\"\n    Tính Average Precision@K cho một user.\n    recommended: list of article_ids (ranked)\n    relevant: set of article_ids (ground truth)\n    \"\"\"\n    if not relevant:\n        return 0.0\n\n    recommended = recommended[:k]\n    hits = 0\n    sum_precisions = 0.0\n\n    for i, item in enumerate(recommended):\n        if item in relevant:\n            hits += 1\n            precision_at_i = hits / (i + 1)\n            sum_precisions += precision_at_i\n\n    # Normalize by min(|relevant|, k)\n    normalizer = min(len(relevant), k)\n    return sum_precisions / normalizer if normalizer > 0 else 0.0\n\n\ndef mean_average_precision_at_k(candidates_dict, ground_truth_dict,\n                                  k=12, verbose=True):\n    \"\"\"\n    Tính MAP@K trên toàn bộ test users.\n    \"\"\"\n    ap_scores = []\n    zero_gt   = 0   # users không có ground truth trong test\n\n    common_users = set(candidates_dict.keys()) & set(ground_truth_dict.keys())\n\n    for cid in common_users:\n        # Ranked list: lấy article_id từ (article_id, score) tuples\n        ranked = [art for art, _ in candidates_dict[cid]]\n        relevant = ground_truth_dict[cid]\n\n        if not relevant:\n            zero_gt += 1\n            continue\n\n        ap = average_precision_at_k(ranked, relevant, k=k)\n        ap_scores.append(ap)\n\n    map_score = np.mean(ap_scores) if ap_scores else 0.0\n\n    if verbose:\n        print(f\"  Evaluated users     : {len(ap_scores):,}\")\n        print(f\"  Users w/ empty GT   : {zero_gt:,}\")\n        print(f\"  AP scores > 0       : {sum(1 for s in ap_scores if s > 0):,} \"\n              f\"({sum(1 for s in ap_scores if s > 0)/len(ap_scores)*100:.1f}%)\")\n        print(f\"  AP distribution:\")\n        ap_arr = np.array(ap_scores)\n        for pct in [25, 50, 75, 90, 95]:\n            print(f\"    P{pct:2d} = {np.percentile(ap_arr, pct):.6f}\")\n\n    return map_score, ap_scores\n\n\nprint(\"=\" * 60)\nprint(\"EVALUATION: MAP@12\")\nprint(\"=\" * 60)\n\nmap12_fpgrowth, ap_scores = mean_average_precision_at_k(\n    candidates, ground_truth, k=12, verbose=True\n)\nprint(f\"\\n MAP@12 (FP-Growth + Popularity): {map12_fpgrowth:.6f}\")","metadata":{"trusted":true},"outputs":[],"execution_count":null},{"cell_type":"code","source":"# ============================================================\n# CELL 34: BASELINE MODELS + COMPARISON TABLE\n#\n# BIỆN LUẬN — Tại sao cần baseline?\n# MAP@12 tuyệt đối không có ý nghĩa nếu không so sánh.\n# Baseline giúp trả lời: \"Recommendation của ta tốt hơn\n# random/popular bao nhiêu?\" — câu hỏi mà giảng viên và\n# business đều muốn nghe con số cụ thể.\n#\n# 3 baselines theo thứ tự phức tạp tăng dần:\n# (1) Random: lower bound — recommendation vô nghĩa\n# (2) Global popularity: standard retail baseline\n# (3) Time-decay popularity: stronger baseline (H&M competition)\n# ============================================================\nimport random as rnd\n\ndef evaluate_random_baseline(ground_truth, article_pool, k=12, n_sample=5000):\n    \"\"\"Random baseline: gợi ý ngẫu nhiên từ article pool.\"\"\"\n    print(\"  Running random baseline...\")\n    rnd.seed(42)\n    test_users = list(ground_truth.keys())[:n_sample]\n    article_list = list(article_pool)\n\n    ap_scores = []\n    for cid in test_users:\n        recommended = rnd.sample(article_list, min(k, len(article_list)))\n        relevant    = ground_truth[cid]\n        ap = average_precision_at_k(recommended, relevant, k=k)\n        ap_scores.append(ap)\n\n    return float(np.mean(ap_scores))\n\n\ndef evaluate_popularity_baseline(ground_truth, pop_df, k=12,\n                                   n_sample=5000, time_decay=False):\n    \"\"\"Popularity baseline (global hoặc time-decay).\"\"\"\n    label = \"time-decay popularity\" if time_decay else \"global popularity\"\n    print(f\"  Running {label} baseline...\")\n\n    top_k_arts = pop_df['article_id'].astype(str).head(k).tolist()\n    test_users = list(ground_truth.keys())[:n_sample]\n\n    ap_scores = []\n    for cid in test_users:\n        relevant = ground_truth[cid]\n        ap = average_precision_at_k(top_k_arts, relevant, k=k)\n        ap_scores.append(ap)\n\n    return float(np.mean(ap_scores))\n\n\ndef evaluate_repurchase_baseline(ground_truth, df_train,\n                                   k=12, n_sample=5000, use_gpu=USE_GPU):\n    \"\"\"Repurchase baseline: gợi ý lại những gì user đã mua.\"\"\"\n    print(\"  Running repurchase baseline...\")\n\n    if use_gpu:\n        rp_pd = df_train[['customer_id', 'article_id']].to_pandas()\n    else:\n        import polars as pl\n        rp_pd = df_train.select(['customer_id', 'article_id']).to_pandas()\n\n    rp_pd['article_id'] = rp_pd['article_id'].astype(str)\n    rp_hist = (\n        rp_pd.groupby('customer_id')['article_id']\n        .apply(list)\n        .to_dict()\n    )\n\n    test_users = list(ground_truth.keys())[:n_sample]\n    ap_scores  = []\n    for cid in test_users:\n        recommended = rp_hist.get(cid, [])[:k]\n        relevant    = ground_truth[cid]\n        ap = average_precision_at_k(recommended, relevant, k=k)\n        ap_scores.append(ap)\n\n    return float(np.mean(ap_scores))\n\n\n# Lấy article pool từ train\nif USE_GPU:\n    all_articles_train = set(df_train['article_id'].astype(str).to_pandas().unique())\nelse:\n    import polars as pl\n    all_articles_train = set(df_train['article_id'].cast(pl.Utf8).unique().to_list())\n\n# ── Chạy tất cả baselines ──\nprint(\"=\" * 60)\nprint(\"BASELINE EVALUATION\")\nprint(\"=\" * 60)\n\nN_EVAL = min(5_000, len(ground_truth))\n\nmap_random   = evaluate_random_baseline(ground_truth, all_articles_train,\n                                         k=12, n_sample=N_EVAL)\nmap_global   = evaluate_popularity_baseline(ground_truth, pop_df,\n                                             k=12, n_sample=N_EVAL,\n                                             time_decay=False)\nmap_decay    = evaluate_popularity_baseline(ground_truth, pop_df,\n                                             k=12, n_sample=N_EVAL,\n                                             time_decay=True)\nmap_repurch  = evaluate_repurchase_baseline(ground_truth, df_train,\n                                             k=12, n_sample=N_EVAL)\n\n# MAP@12 FP-Growth đã tính ở Cell 33\nmap_fp = map12_fpgrowth\n\n\n# ── Bảng so sánh ──\ncomparison = {\n    'Model'                  : ['Random',\n                                'Global popularity',\n                                'Time-decay popularity',\n                                'Repurchase history',\n                                'FP-Growth + Popularity (ours)'],\n    'MAP@12'                 : [map_random, map_global, map_decay,\n                                map_repurch, map_fp],\n    'vs Random (+%)'         : [0.0,\n                                (map_global  - map_random) / (map_random + 1e-9) * 100,\n                                (map_decay   - map_random) / (map_random + 1e-9) * 100,\n                                (map_repurch - map_random) / (map_random + 1e-9) * 100,\n                                (map_fp      - map_random) / (map_random + 1e-9) * 100],\n    'Approach'               : ['Lower bound',\n                                'Non-personalized',\n                                'Non-personalized + recency',\n                                'User-level history',\n                                'Association rules + recency'],\n}\n\ncomp_df = pd.DataFrame(comparison)\ncomp_df['MAP@12']       = comp_df['MAP@12'].map('{:.6f}'.format)\ncomp_df['vs Random (+%)'] = comp_df['vs Random (+%)'].map('{:+.1f}%'.format)\n\nprint(\"\\n\" + \"=\" * 70)\nprint(\"  COMPARISON TABLE — MAP@12 (test = last 7 days)\")\nprint(\"=\" * 70)\nprint(comp_df.to_string(index=False))\nprint(\"=\" * 70)\n\nprint(f\"\\n  NHẬN XÉT:\")\nprint(f\"  - FP-Growth + Popularity đạt MAP@12 = {map_fp:.6f}\")\nprint(f\"  - Cải thiện {(map_fp-map_decay)/(map_decay+1e-9)*100:+.1f}% so với time-decay baseline\")\nprint(f\"  - Association rules đóng góp personalization layer\")\nprint(f\"    mà pure popularity không có được\")","metadata":{"trusted":true},"outputs":[],"execution_count":null},{"cell_type":"code","source":"# ============================================================\n# CELL 34: BASELINE MODELS + COMPARISON TABLE\n#\n# BIỆN LUẬN — Tại sao cần baseline?\n# MAP@12 tuyệt đối không có ý nghĩa nếu không so sánh.\n# Baseline giúp trả lời: \"Recommendation của ta tốt hơn\n# random/popular bao nhiêu?\" — câu hỏi mà giảng viên và\n# business đều muốn nghe con số cụ thể.\n#\n# 3 baselines theo thứ tự phức tạp tăng dần:\n# (1) Random: lower bound — recommendation vô nghĩa\n# (2) Global popularity: standard retail baseline\n# (3) Time-decay popularity: stronger baseline (H&M competition)\n# ============================================================\nimport random as rnd\n\ndef evaluate_random_baseline(ground_truth, article_pool, k=12, n_sample=5000):\n    \"\"\"Random baseline: gợi ý ngẫu nhiên từ article pool.\"\"\"\n    print(\"  Running random baseline...\")\n    rnd.seed(42)\n    test_users = list(ground_truth.keys())[:n_sample]\n    article_list = list(article_pool)\n\n    ap_scores = []\n    for cid in test_users:\n        recommended = rnd.sample(article_list, min(k, len(article_list)))\n        relevant    = ground_truth[cid]\n        ap = average_precision_at_k(recommended, relevant, k=k)\n        ap_scores.append(ap)\n\n    return float(np.mean(ap_scores))\n\n\ndef evaluate_popularity_baseline(ground_truth, pop_df, k=12,\n                                   n_sample=5000, time_decay=False):\n    \"\"\"Popularity baseline (global hoặc time-decay).\"\"\"\n    label = \"time-decay popularity\" if time_decay else \"global popularity\"\n    print(f\"  Running {label} baseline...\")\n\n    top_k_arts = pop_df['article_id'].astype(str).head(k).tolist()\n    test_users = list(ground_truth.keys())[:n_sample]\n\n    ap_scores = []\n    for cid in test_users:\n        relevant = ground_truth[cid]\n        ap = average_precision_at_k(top_k_arts, relevant, k=k)\n        ap_scores.append(ap)\n\n    return float(np.mean(ap_scores))\n\n\ndef evaluate_repurchase_baseline(ground_truth, df_train,\n                                   k=12, n_sample=5000, use_gpu=USE_GPU):\n    \"\"\"Repurchase baseline: gợi ý lại những gì user đã mua.\"\"\"\n    print(\"  Running repurchase baseline...\")\n\n    if use_gpu:\n        rp_pd = df_train[['customer_id', 'article_id']].to_pandas()\n    else:\n        import polars as pl\n        rp_pd = df_train.select(['customer_id', 'article_id']).to_pandas()\n\n    rp_pd['article_id'] = rp_pd['article_id'].astype(str)\n    rp_hist = (\n        rp_pd.groupby('customer_id')['article_id']\n        .apply(list)\n        .to_dict()\n    )\n\n    test_users = list(ground_truth.keys())[:n_sample]\n    ap_scores  = []\n    for cid in test_users:\n        recommended = rp_hist.get(cid, [])[:k]\n        relevant    = ground_truth[cid]\n        ap = average_precision_at_k(recommended, relevant, k=k)\n        ap_scores.append(ap)\n\n    return float(np.mean(ap_scores))\n\n\n# Lấy article pool từ train\nif USE_GPU:\n    all_articles_train = set(df_train['article_id'].astype(str).to_pandas().unique())\nelse:\n    import polars as pl\n    all_articles_train = set(df_train['article_id'].cast(pl.Utf8).unique().to_list())\n\n# ── Chạy tất cả baselines ──\nprint(\"=\" * 60)\nprint(\"BASELINE EVALUATION\")\nprint(\"=\" * 60)\n\nN_EVAL = min(5_000, len(ground_truth))\n\nmap_random   = evaluate_random_baseline(ground_truth, all_articles_train,\n                                         k=12, n_sample=N_EVAL)\nmap_global   = evaluate_popularity_baseline(ground_truth, pop_df,\n                                             k=12, n_sample=N_EVAL,\n                                             time_decay=False)\nmap_decay    = evaluate_popularity_baseline(ground_truth, pop_df,\n                                             k=12, n_sample=N_EVAL,\n                                             time_decay=True)\nmap_repurch  = evaluate_repurchase_baseline(ground_truth, df_train,\n                                             k=12, n_sample=N_EVAL)\n\n# MAP@12 FP-Growth đã tính ở Cell 33\nmap_fp = map12_fpgrowth\n\n\n# ── Bảng so sánh ──\ncomparison = {\n    'Model'                  : ['Random',\n                                'Global popularity',\n                                'Time-decay popularity',\n                                'Repurchase history',\n                                'FP-Growth + Popularity (ours)'],\n    'MAP@12'                 : [map_random, map_global, map_decay,\n                                map_repurch, map_fp],\n    'vs Random (+%)'         : [0.0,\n                                (map_global  - map_random) / (map_random + 1e-9) * 100,\n                                (map_decay   - map_random) / (map_random + 1e-9) * 100,\n                                (map_repurch - map_random) / (map_random + 1e-9) * 100,\n                                (map_fp      - map_random) / (map_random + 1e-9) * 100],\n    'Approach'               : ['Lower bound',\n                                'Non-personalized',\n                                'Non-personalized + recency',\n                                'User-level history',\n                                'Association rules + recency'],\n}\n\ncomp_df = pd.DataFrame(comparison)\ncomp_df['MAP@12']       = comp_df['MAP@12'].map('{:.6f}'.format)\ncomp_df['vs Random (+%)'] = comp_df['vs Random (+%)'].map('{:+.1f}%'.format)\n\nprint(\"\\n\" + \"=\" * 70)\nprint(\"  COMPARISON TABLE — MAP@12 (test = last 7 days)\")\nprint(\"=\" * 70)\nprint(comp_df.to_string(index=False))\nprint(\"=\" * 70)\n\nprint(f\"\\n  NHẬN XÉT:\")\nprint(f\"  - FP-Growth + Popularity đạt MAP@12 = {map_fp:.6f}\")\nprint(f\"  - Cải thiện {(map_fp-map_decay)/(map_decay+1e-9)*100:+.1f}% so với time-decay baseline\")\nprint(f\"  - Association rules đóng góp personalization layer\")\nprint(f\"    mà pure popularity không có được\")","metadata":{"trusted":true},"outputs":[],"execution_count":null},{"cell_type":"code","source":"# ============================================================\n# CELL 36: LƯU KẾT QUẢ EVALUATION → MongoDB (audit trail)\n# ============================================================\nfrom datetime import datetime, timezone\n\neval_doc = {\n    \"run_id\"       : f\"eval-{datetime.now().strftime('%Y%m%d-%H%M')}\",\n    \"cutoff_date\"  : cutoff_str,\n    \"test_users\"   : N_EVAL,\n    \"metric\"       : \"MAP@12\",\n    \"results\"      : {\n        \"random\"         : round(map_random,  6),\n        \"global_pop\"     : round(map_global,  6),\n        \"time_decay_pop\" : round(map_decay,   6),\n        \"repurchase\"     : round(map_repurch, 6),\n        \"fp_growth\"      : round(map_fp,      6),\n    },\n    \"fp_growth_config\" : {\n        \"min_support\"    : BEST_SUPPORT,\n        \"min_confidence\" : BEST_CONFIDENCE,\n        \"n_rules\"        : len(fp_rules),\n        \"basket_sample\"  : 80_000,\n        \"max_len\"        : 3,\n    },\n    \"created_at\" : datetime.now(timezone.utc),\n}\n\ndb[\"evaluation_runs\"].insert_one(eval_doc)\nprint(f\"   Evaluation results saved to MongoDB: evaluation_runs\")\nprint(f\"   run_id: {eval_doc['run_id']}\")\nprint(f\"   MAP@12 (FP-Growth): {eval_doc['results']['fp_growth']}\")","metadata":{"trusted":true},"outputs":[],"execution_count":null},{"cell_type":"code","source":"","metadata":{"trusted":true},"outputs":[],"execution_count":null}]}