{"metadata":{"kernelspec":{"language":"python","display_name":"Python 3","name":"python3"},"language_info":{"name":"python","version":"3.11.13","mimetype":"text/x-python","codemirror_mode":{"name":"ipython","version":3},"pygments_lexer":"ipython3","nbconvert_exporter":"python","file_extension":".py"},"kaggle":{"accelerator":"tpuV5e8","dataSources":[{"sourceId":31254,"databundleVersionId":3103714,"sourceType":"competition"},{"sourceId":14112508,"sourceType":"datasetVersion","datasetId":8989713},{"sourceId":14241442,"sourceType":"datasetVersion","datasetId":9085906},{"sourceId":14258748,"sourceType":"datasetVersion","datasetId":9098345}],"dockerImageVersionId":31192,"isInternetEnabled":true,"language":"python","sourceType":"notebook","isGpuEnabled":false}},"nbformat_minor":4,"nbformat":4,"cells":[{"cell_type":"code","source":"import pandas as pd\nimport numpy as np\nimport os\nimport gc\nimport re\nfrom datetime import timedelta\nimport traceback\nimport pyarrow as pa\nimport pyarrow.parquet as pq\n","metadata":{"trusted":true,"execution":{"iopub.status.busy":"2025-12-22T16:12:04.112247Z","iopub.execute_input":"2025-12-22T16:12:04.113483Z","iopub.status.idle":"2025-12-22T16:12:05.440712Z","shell.execute_reply.started":"2025-12-22T16:12:04.113446Z","shell.execute_reply":"2025-12-22T16:12:05.439187Z"}},"outputs":[],"execution_count":null},{"cell_type":"code","source":"# --- CẤU HÌNH CONFIG ---\nclass Config:\n    # 1. Đường dẫn dataset gốc\n    RAW_DATA_DIR = '/kaggle/input/h-and-m-personalized-fashion-recommendations/'\n    \n    # 2. Đường dẫn dataset Candidate\n    CANDIDATE_DIR = '/kaggle/input/h-and-m-recall-model-dataset/'\n    \n    # 3. Thư mục xuất file kết quả\n    OUTPUT_DIR = '/kaggle/working/enriched_data/'\n    \n    # 4. Các Feature CẦN LOẠI BỎ\n    DROP_FEATURES = ['dssm_similarity', 'yt_similarity', 'wv_similarity', 'label']\n    \n    # 5. Kích thước chunk\n    CHUNK_SIZE = 200_000 # Giảm xuống 200k cho an toàn RAM\n","metadata":{"trusted":true,"execution":{"iopub.status.busy":"2025-12-22T16:12:05.443024Z","iopub.execute_input":"2025-12-22T16:12:05.443543Z","iopub.status.idle":"2025-12-22T16:19:46.199525Z","shell.execute_reply.started":"2025-12-22T16:12:05.443514Z","shell.execute_reply":"2025-12-22T16:19:46.198536Z"}},"outputs":[],"execution_count":null},{"cell_type":"code","source":"# =============================================================================\n# 1. UTILS\n# =============================================================================\ndef reduce_mem_usage(df):\n    \"\"\"Giảm dung lượng RAM, bỏ qua cột datetime\"\"\"\n    for col in df.columns:\n        col_type = df[col].dtype\n        if col_type != object and not np.issubdtype(col_type, np.datetime64):\n            c_min = df[col].min()\n            c_max = df[col].max()\n            if str(col_type)[:3] == 'int':\n                if c_min > np.iinfo(np.int32).min and c_max < np.iinfo(np.int32).max:\n                    df[col] = df[col].astype(np.int32)\n                else:\n                    df[col] = df[col].astype(np.int64)\n            else:\n                if c_min > np.finfo(np.float32).min and c_max < np.finfo(np.float32).max:\n                    df[col] = df[col].astype(np.float32)\n                else:\n                    df[col] = df[col].astype(np.float32)\n    return df\n","metadata":{"trusted":true,"execution":{"iopub.status.busy":"2025-12-22T16:19:46.200831Z","iopub.execute_input":"2025-12-22T16:19:46.201205Z","iopub.status.idle":"2025-12-22T16:19:47.455174Z","shell.execute_reply.started":"2025-12-22T16:19:46.201167Z","shell.execute_reply":"2025-12-22T16:19:47.454171Z"}},"outputs":[],"execution_count":null},{"cell_type":"code","source":"# =============================================================================\n# 2. CORE CLASS\n# =============================================================================\nclass HMFeaturePipeline:\n    def __init__(self):\n        self.raw_path = Config.RAW_DATA_DIR\n        self.pqt_path = Config.CANDIDATE_DIR\n        self.output_dir = Config.OUTPUT_DIR\n        \n        if not os.path.exists(self.output_dir): \n            os.makedirs(self.output_dir)\n            \n    def prepare_resources(self):\n        print(\"\\n=== [1] PREPARING RESOURCES ===\")\n        \n        # --- A. LOAD CUSTOMERS ---\n        print(\"-> Loading Customers...\")\n        customers = pd.read_csv(self.raw_path + 'customers.csv')\n        self.cust_id_map = dict(zip(customers['customer_id'], customers.index))\n        \n        customers['customer_id'] = customers.index.astype('int32')\n        customers['age'] = customers['age'].fillna(customers['age'].mean()).astype(np.int8)\n        self.customers_df = reduce_mem_usage(customers[['customer_id', 'age']])\n        del customers; gc.collect()\n\n        # --- B. LOAD ARTICLES ---\n        print(\"-> Loading Articles...\")\n        articles = pd.read_csv(self.raw_path + 'articles.csv', dtype={'article_id': str})\n        self.article_id_map = dict(zip(articles['article_id'], articles.index))\n        del articles; gc.collect()\n        \n        # --- C. LOAD TRANSACTIONS ---\n        print(\"-> Loading Transactions (Last 5 weeks)...\")\n        df_trans = pd.read_csv(self.raw_path + 'transactions_train.csv', \n                               dtype={'article_id': str, 'sales_channel_id': 'int8'},\n                               parse_dates=['t_dat'])\n        \n        print(\"   Mapping IDs...\")\n        df_trans['customer_id'] = df_trans['customer_id'].map(self.cust_id_map).fillna(-1).astype('int32')\n        df_trans['article_id'] = df_trans['article_id'].map(self.article_id_map).fillna(-1).astype('int32')\n        \n        self.max_date = df_trans['t_dat'].max()\n        start_date = self.max_date - timedelta(days=35) \n        self.trans_df = df_trans[df_trans['t_dat'] >= start_date].copy()\n        self.trans_df['week'] = (self.max_date - self.trans_df['t_dat']).dt.days // 7\n        self.trans_df = reduce_mem_usage(self.trans_df)\n        del df_trans; gc.collect()\n        \n        # --- D. CALCULATE HELPERS ---\n        print(\"-> Calculating Helpers...\")\n        self.user_avg_spend = self.trans_df.groupby('customer_id')['price'].mean().reset_index(name='user_avg_price')\n        self.user_avg_spend = reduce_mem_usage(self.user_avg_spend)\n        \n        self.item_curr_price = self.trans_df.groupby('article_id')['price'].mean().reset_index(name='item_curr_price')\n        self.item_curr_price = reduce_mem_usage(self.item_curr_price)\n        \n        merged = self.trans_df.merge(self.customers_df, on='customer_id', how='left')\n        self.item_target_age = merged.groupby('article_id')['age'].mean().reset_index(name='item_target_age')\n        self.item_target_age = reduce_mem_usage(self.item_target_age)\n        \n        self.item_first_sale = self.trans_df.groupby('article_id')['t_dat'].min().reset_index(name='first_sale_date')\n        \n        print(\"   Building Sales Trend Lookup...\")\n        self.sales_lookup = self.trans_df.groupby(['week', 'article_id']).size().to_dict()\n        \n        del merged; gc.collect()\n        print(\"-> Resources Ready.\")\n\n    def process_file(self, filename):\n        match = re.search(r'week(\\d+)', filename)\n        if not match: return\n        \n        week_num = int(match.group(1))\n        input_path = os.path.join(self.pqt_path, filename)\n        \n        base_name = os.path.splitext(filename)[0]\n        output_file = f\"enriched_{base_name}.pqt\"\n        output_path = os.path.join(self.output_dir, output_file)\n        \n        if os.path.exists(output_path): os.remove(output_path)\n        \n        print(f\"\\n>> Processing: {filename} (Week {week_num})\")\n        \n        try:\n            # Load Label nếu có\n            lbl_df = None\n            label_path = f\"{self.pqt_path}/week{week_num}_label.pqt\"\n            if os.path.exists(label_path):\n                lbl = pd.read_parquet(label_path)\n                lbl = lbl.explode('article_id')\n                lbl['article_id'] = lbl['article_id'].astype('int32')\n                lbl['customer_id'] = lbl['customer_id'].astype('int32')\n                lbl['target'] = 1\n                lbl_df = lbl[['customer_id', 'article_id', 'target']]\n                del lbl\n            \n            # --- CHUNKING VỚI PYARROW (FIXED) ---\n            # Sử dụng ParquetFile và iter_batches thay vì pd.read_parquet(chunksize)\n            pf = pq.ParquetFile(input_path)\n            writer = None\n            \n            # Loop qua từng batch (chunk)\n            for i, batch in enumerate(pf.iter_batches(batch_size=Config.CHUNK_SIZE)):\n                df = batch.to_pandas()\n                \n                # 1. CLEANING & CASTING\n                df.drop(columns=[c for c in Config.DROP_FEATURES if c in df.columns], inplace=True)\n                df['customer_id'] = df['customer_id'].astype('int32')\n                df['article_id'] = df['article_id'].astype('int32')\n                \n                # 2. MERGE TARGET\n                if lbl_df is not None:\n                    df = df.merge(lbl_df, on=['customer_id', 'article_id'], how='left')\n                    df['target'] = df['target'].fillna(0).astype(np.int8)\n                else:\n                    df['target'] = np.int8(0)\n                \n                # 3. ENRICHMENT\n                if 'age' in df.columns: df.drop(columns=['age'], inplace=True)\n                \n                df = df.merge(self.customers_df, on='customer_id', how='left')\n                df = df.merge(self.user_avg_spend, on='customer_id', how='left')\n                df['user_avg_price'] = df['user_avg_price'].fillna(0.03).astype(np.float32)\n                \n                df = df.merge(self.item_curr_price, on='article_id', how='left')\n                df = df.merge(self.item_target_age, on='article_id', how='left')\n                df = df.merge(self.item_first_sale, on='article_id', how='left')\n                \n                # 4. CALCULATE FEATURES\n                df['price_sensitivity'] = (df['item_curr_price'] - df['user_avg_price']).astype(np.float32)\n                df['price_sensitivity'] = df['price_sensitivity'].fillna(0.0)\n                \n                df['age_diff'] = abs(df['age'] - df['item_target_age'])\n                df['age_diff'] = df['age_diff'].fillna(10.0).astype(np.float32)\n                \n                curr_date = self.max_date - timedelta(days=week_num*7)\n                df['days_since_release'] = (curr_date - df['first_sale_date']).dt.days\n                df['is_new_arrival'] = (df['days_since_release'] <= 30).astype(np.int8)\n                df['is_new_arrival'] = df['is_new_arrival'].fillna(0)\n                \n                # Sales Trend (Map Dict)\n                prev_w = week_num + 1\n                prev2_w = week_num + 2\n                \n                df['sales_w1'] = df['article_id'].map(lambda x: self.sales_lookup.get((prev_w, x), 0)).astype(np.int16)\n                df['sales_w2'] = df['article_id'].map(lambda x: self.sales_lookup.get((prev2_w, x), 0)).astype(np.int16)\n                df['sales_trend'] = ((df['sales_w1'] - df['sales_w2']) / (df['sales_w2'] + 1)).astype(np.float32)\n                \n                # 5. CLEANUP & SAVE\n                cols_drop = ['first_sale_date', 'item_curr_price', 'item_target_age', 'user_avg_price', 'sales_w1', 'sales_w2']\n                df.drop(columns=cols_drop, inplace=True, errors='ignore')\n                \n                # Ghi vào file Parquet Output\n                table = pa.Table.from_pandas(df)\n                if writer is None:\n                    writer = pq.ParquetWriter(output_path, table.schema)\n                writer.write_table(table)\n                \n                del df, table; gc.collect()\n                \n            if writer: writer.close()\n            print(f\"   -> Finished. Saved to: {output_file}\")\n            \n            del lbl_df; gc.collect()\n            \n        except Exception as e:\n            print(f\"!!! Error processing {filename}: {e}\")\n            traceback.print_exc()\n\n    def run_all(self):\n        print(\"\\n--- [2] STARTING BATCH PROCESSING ---\")\n        files = [f for f in os.listdir(self.pqt_path) if 'candidate' in f and f.endswith('.pqt')]\n        files.sort()\n        \n        print(f\"Found {len(files)} files to process.\")\n        for f in files:\n            self.process_file(f)\n        print(\"\\n=== ALL DONE ===\")\n","metadata":{"trusted":true,"execution":{"iopub.status.busy":"2025-12-22T16:19:47.457454Z","iopub.execute_input":"2025-12-22T16:19:47.457754Z","iopub.status.idle":"2025-12-22T16:19:47.478644Z","shell.execute_reply.started":"2025-12-22T16:19:47.457731Z","shell.execute_reply":"2025-12-22T16:19:47.477652Z"}},"outputs":[],"execution_count":null},{"cell_type":"code","source":"if __name__ == \"__main__\":\n    pipeline = HMFeaturePipeline()\n    pipeline.prepare_resources()\n    pipeline.run_all()","metadata":{"trusted":true},"outputs":[],"execution_count":null},{"cell_type":"code","source":"import pandas as pd\nimport numpy as np\nimport os\nimport matplotlib.pyplot as plt\nimport seaborn as sns\n\n# Cấu hình\npd.set_option('display.max_columns', None)\nsns.set_style(\"whitegrid\")\n\n# Đường dẫn Output từ bước trước\nOUTPUT_DIR = '/kaggle/working/enriched_data/'\n\ndef inspect_enriched_files():\n    print(f\"=== KIỂM TRA THƯ MỤC: {OUTPUT_DIR} ===\\n\")\n    \n    if not os.path.exists(OUTPUT_DIR):\n        print(\"❌ Thư mục không tồn tại. Bạn đã chạy pipeline chưa?\")\n        return\n\n    files = sorted([f for f in os.listdir(OUTPUT_DIR) if f.endswith('.pqt')])\n    \n    if not files:\n        print(\"❌ Thư mục rỗng. Chưa có file nào được tạo.\")\n        return\n        \n    print(f\"✅ Tìm thấy {len(files)} file: {files}\\n\")\n    \n    # --- CHỌN FILE ĐẦU TIÊN ĐỂ SOI ---\n    sample_file = files[0] # Hoặc chọn file cụ thể: 'enriched_week1_candidate.pqt'\n    file_path = os.path.join(OUTPUT_DIR, sample_file)\n    \n    print(f\"🔎 ĐANG PHÂN TÍCH FILE MẪU: {sample_file}\")\n    print(\"-\" * 50)\n    \n    # Load data\n    df = pd.read_parquet(file_path)\n    \n    # 1. Overview\n    print(f\"1. Kích thước (Shape): {df.shape}\")\n    print(f\"2. Danh sách cột ({len(df.columns)} cols):\")\n    print(df.columns.tolist())\n    \n    # 2. Kiểm tra các Feature Mới (Quan trọng nhất)\n    new_features = ['price_sensitivity', 'age_diff', 'is_new_arrival', 'sales_trend']\n    print(\"\\n3. Thống kê các Feature Mới (Kiểm tra xem có NaN không):\")\n    display(df[new_features].describe())\n    \n    # 3. Xem mẫu dữ liệu\n    print(\"\\n4. Dữ liệu mẫu (Head):\")\n    display(df.head(5))\n    \n    # 4. Kiểm tra Target\n    print(\"\\n5. Phân phối Target (Label):\")\n    print(df['target'].value_counts(normalize=True))\n    \n    # 5. Trực quan hóa nhanh\n    print(\"\\n6. Biểu đồ phân phối Feature mới:\")\n    plt.figure(figsize=(15, 4))\n    \n    plt.subplot(1, 3, 1)\n    sns.histplot(df['price_sensitivity'].sample(10000), bins=30, kde=True, color='green')\n    plt.title('Price Sensitivity (Âm = Rẻ hơn thói quen)')\n    \n    plt.subplot(1, 3, 2)\n    sns.histplot(df['age_diff'].sample(10000), bins=30, kde=True, color='orange')\n    plt.title('Age Diff (Càng nhỏ càng hợp tuổi)')\n    \n    plt.subplot(1, 3, 3)\n    sns.histplot(df['sales_trend'].sample(10000), bins=30, kde=True, color='purple')\n    plt.title('Sales Trend (>0 là đang Hot)')\n    \n    plt.tight_layout()\n    plt.show()\n\n# Chạy hàm kiểm tra\ninspect_enriched_files()","metadata":{"trusted":true},"outputs":[],"execution_count":null}]}