{"metadata":{"kernelspec":{"language":"python","display_name":"Python 3","name":"python3"},"language_info":{"pygments_lexer":"ipython3","nbconvert_exporter":"python","version":"3.6.4","file_extension":".py","codemirror_mode":{"name":"ipython","version":3},"name":"python","mimetype":"text/x-python"},"kaggle":{"accelerator":"tpuV5e8","dataSources":[{"sourceId":31254,"databundleVersionId":3103714,"sourceType":"competition"},{"sourceId":14112508,"sourceType":"datasetVersion","datasetId":8989713}],"dockerImageVersionId":31192,"isInternetEnabled":true,"language":"python","sourceType":"notebook","isGpuEnabled":false}},"nbformat_minor":4,"nbformat":4,"cells":[{"cell_type":"code","source":"# =============================================================================\n# === H&M PIPELINE: CHUNK PROCESSING -> SINGLE PARQUET FILE OUTPUT ===\n# =============================================================================\n\nimport 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\npd.set_option('display.max_columns', None)\n\n# --- 1. HÀM GIẢM BỘ NHỚ ---\ndef reduce_mem_usage(df):\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\nclass HMSingleFileEnricher:\n    def __init__(self, raw_path, pqt_path, output_dir='/kaggle/working/'):\n        self.raw_path = raw_path\n        self.pqt_path = pqt_path\n        self.output_dir = output_dir\n        if not os.path.exists(output_dir): os.makedirs(output_dir)\n        \n    def load_raw_resources(self):\n        print(\"--- [1] Loading Raw Data & Creating ID Maps ---\")\n        \n        # 1. CUSTOMERS\n        print(\"-> Loading Customers...\")\n        self.customers = pd.read_csv(self.raw_path + 'customers.csv')\n        cust_id_to_idx = dict(zip(self.customers['customer_id'], self.customers.index))\n        \n        self.customers['customer_id'] = self.customers.index.astype('int32')\n        self.customers['age'] = self.customers['age'].fillna(self.customers['age'].mean()).astype(np.int8)\n        self.customers = self.customers[['customer_id', 'age']]\n        self.customers = reduce_mem_usage(self.customers)\n\n        # 2. ARTICLES\n        print(\"-> Loading Articles...\")\n        self.articles = pd.read_csv(self.raw_path + 'articles.csv', dtype={'article_id': str})\n        article_id_to_idx = dict(zip(self.articles['article_id'], self.articles.index))\n        \n        # 3. TRANSACTIONS\n        print(\"-> Loading Transactions...\")\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(cust_id_to_idx).fillna(-1).astype('int32')\n        df_trans['article_id'] = df_trans['article_id'].map(article_id_to_idx).fillna(-1).astype('int32')\n        \n        self.max_date = df_trans['t_dat'].max()\n        start_date = self.max_date - timedelta(days=40)\n        self.trans = df_trans[df_trans['t_dat'] >= start_date].copy()\n        self.trans['week'] = (self.max_date - self.trans['t_dat']).dt.days // 7\n        \n        del df_trans, cust_id_to_idx, article_id_to_idx; gc.collect()\n        \n        # --- CALCULATE HELPERS ---\n        print(\"-> Calculating Helpers...\")\n        \n        # A. User Stats\n        self.user_avg_spend = self.trans.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        # B. Item Target Age\n        trans_with_age = self.trans.merge(self.customers[['customer_id', 'age']], on='customer_id')\n        self.item_target_age = trans_with_age.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        # C. Item Price\n        self.item_curr_price = self.trans.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        # D. First Sale\n        self.item_first_sale = self.trans.groupby('article_id')['t_dat'].min().reset_index(name='first_sale_date')\n        \n        # E. Sales Lookup Dict\n        print(\"-> Creating Sales Lookup Dict...\")\n        self.sales_lookup = self.trans.groupby(['week', 'article_id']).size().to_dict()\n        \n        del trans_with_age; gc.collect()\n        print(\"-> Helpers Ready.\")\n\n    def process_and_append(self, filename, chunk_size=500000):\n        # Xác định tuần\n        match = re.search(r'week(\\d+)', filename)\n        if not match: return\n        week_num = int(match.group(1))\n        \n        file_path = os.path.join(self.pqt_path, filename)\n        print(f\"\\n>> Processing File: {filename} (Week {week_num})...\")\n        \n        try:\n            df_full = pd.read_parquet(file_path)\n            \n            # Chuẩn bị Label\n            label_path = f\"{self.pqt_path}/week{week_num}_label.pqt\"\n            lbl_df = None\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            # --- CẤU HÌNH OUTPUT SINGLE FILE ---\n            # Nếu tên file gốc là week0_candidate_0.pqt -> Output là enriched_week0_candidate_0.pqt\n            # Nếu muốn gộp chung hết week0, logic sẽ phức tạp hơn (cần check file tồn tại).\n            # Ở đây ta giữ nguyên mapping 1 input file -> 1 output enriched file.\n            base_name = os.path.splitext(filename)[0]\n            output_file = f\"enriched_{base_name}.pqt\"\n            save_path = os.path.join(self.output_dir, output_file)\n            \n            # Xóa file cũ nếu tồn tại\n            if os.path.exists(save_path): os.remove(save_path)\n            \n            # Khởi tạo Parquet Writer\n            pqwriter = None\n            \n            total_rows = len(df_full)\n            num_chunks = (total_rows // chunk_size) + 1\n            print(f\"   Writing to: {output_file} ({num_chunks} chunks)\")\n            \n            for i in range(num_chunks):\n                start = i * chunk_size\n                end = min((i + 1) * chunk_size, total_rows)\n                if start >= end: break\n                \n                # --- PROCESS CHUNK ---\n                df = df_full.iloc[start:end].copy()\n                \n                # Ép kiểu ID\n                df['customer_id'] = df['customer_id'].astype('int32')\n                df['article_id'] = df['article_id'].astype('int32')\n                \n                # Drop rác\n                drop_cols = ['dssm_similarity', 'yt_similarity', 'wv_similarity', 'label']\n                df.drop(columns=[c for c in drop_cols if c in df.columns], inplace=True)\n                \n                # 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                # Merge Features\n                if 'age' in df.columns: df.drop(columns=['age'], inplace=True)\n                df = df.merge(self.customers, 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                # Calc 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 (Lookup)\n                prev_week = week_num + 1\n                prev2_week = week_num + 2\n                df['sales_w1'] = df['article_id'].map(lambda x: self.sales_lookup.get((prev_week, x), 0)).astype(np.int16)\n                df['sales_w2'] = df['article_id'].map(lambda x: self.sales_lookup.get((prev2_week, x), 0)).astype(np.int16)\n                df['sales_trend'] = ((df['sales_w1'] - df['sales_w2']) / (df['sales_w2'] + 1)).astype(np.float32)\n                \n                # Cleanup\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                # --- WRITE TO PARQUET (APPEND MODE) ---\n                table = pa.Table.from_pandas(df)\n                \n                if pqwriter is None:\n                    # Tạo writer lần đầu tiên\n                    pqwriter = pq.ParquetWriter(save_path, table.schema)\n                \n                pqwriter.write_table(table)\n                \n                del df, table; gc.collect()\n            \n            # Đóng writer để hoàn tất file\n            if pqwriter:\n                pqwriter.close()\n                \n            print(f\"   -> Finished: {output_file}\")\n            del df_full, lbl_df; gc.collect()\n            \n        except Exception as e:\n            print(f\"!!! Error processing {filename}: {e}\")\n            print(traceback.format_exc())\n\n    def process_all_files(self):\n        print(\"\\n--- [2] Scanning & Processing All Files ---\")\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.\")\n        for f in files:\n            self.process_and_append(f)\n            \n        print(\"\\n=== ALL FILES PROCESSED SUCCESSFULLY ===\")\n\n# =============================================================================\nif __name__ == \"__main__\":\n    RAW_PATH = '/kaggle/input/h-and-m-personalized-fashion-recommendations/'\n    PQT_PATH = '/kaggle/input/h-and-m-recall-model-dataset/' \n    \n    enricher = HMSingleFileEnricher(RAW_PATH, PQT_PATH)\n    enricher.load_raw_resources()\n    enricher.process_all_files()","metadata":{"trusted":true},"outputs":[],"execution_count":null}]}