{"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":"datasetVersion","sourceId":15869438,"datasetId":9951844,"databundleVersionId":16822045},{"sourceType":"datasetVersion","sourceId":15880486,"datasetId":9869115,"databundleVersionId":16833966}],"dockerImageVersionId":31329,"isInternetEnabled":true,"language":"python","sourceType":"notebook","isGpuEnabled":true}},"nbformat_minor":4,"nbformat":4,"cells":[{"cell_type":"code","source":"# This Python 3 environment comes with many helpful analytics libraries installed\n# It is defined by the kaggle/python Docker image: https://github.com/kaggle/docker-python\n# For example, here's several helpful packages to load\n\nimport numpy as np # linear algebra\nimport pandas as pd # data processing, CSV file I/O (e.g. pd.read_csv)\n\n# Input data files are available in the read-only \"../input/\" directory\n# For example, running this (by clicking run or pressing Shift+Enter) will list all files under the input directory\n\nimport os\nfor dirname, _, filenames in os.walk('/kaggle/input'):\n    for filename in filenames:\n        print(os.path.join(dirname, filename))\n\n# You can write up to 20GB to the current directory (/kaggle/working/) that gets preserved as output when you create a version using \"Save & Run All\" \n# You can also write temporary files to /kaggle/temp/, but they won't be saved outside of the current session","metadata":{"_uuid":"8f2839f25d086af736a60e9eeb907d3b93b6e0e5","_cell_guid":"b1076dfc-b9ad-4769-8c92-a6c4dae69d19","trusted":true,"execution":{"iopub.status.busy":"2026-04-22T04:37:04.158358Z","iopub.execute_input":"2026-04-22T04:37:04.158627Z","iopub.status.idle":"2026-04-22T04:37:05.460933Z","shell.execute_reply.started":"2026-04-22T04:37:04.158592Z","shell.execute_reply":"2026-04-22T04:37:05.460071Z"},"scrolled":true},"outputs":[],"execution_count":null},{"cell_type":"code","source":"!pip install openai-clip --quiet","metadata":{"trusted":true,"execution":{"iopub.status.busy":"2026-04-22T04:37:05.461982Z","iopub.execute_input":"2026-04-22T04:37:05.463437Z","iopub.status.idle":"2026-04-22T04:37:12.391693Z","shell.execute_reply.started":"2026-04-22T04:37:05.463410Z","shell.execute_reply":"2026-04-22T04:37:12.390768Z"}},"outputs":[],"execution_count":null},{"cell_type":"code","source":"import os\nimport gc\nimport numpy as np\nimport pandas as pd\nfrom collections import defaultdict\n \nfrom sklearn.preprocessing import OneHotEncoder, normalize\nfrom sklearn.decomposition import TruncatedSVD\nfrom sklearn.cluster import KMeans\n\nimport matplotlib.pyplot as plt\nimport seaborn as sns\nfrom sentence_transformers import SentenceTransformer\nimport torch\nfrom clip import load\nmodel_clip, preprocess = load('ViT-B/32', device='cuda')\nfrom PIL import Image\nfrom torch.utils.data import Dataset, DataLoader\nimport os\nfrom sklearn.preprocessing import normalize\nfrom sklearn.decomposition import PCA\nfrom sklearn.metrics import silhouette_score\nfrom scipy.sparse import csr_matrix\n\nimport polars as pl\nimport pandas as pd\nimport numpy as np\nimport gc\nimport torch\nimport pyarrow.parquet as pq\nimport pyarrow as pa","metadata":{"trusted":true,"execution":{"iopub.status.busy":"2026-04-22T04:37:12.394257Z","iopub.execute_input":"2026-04-22T04:37:12.394528Z","iopub.status.idle":"2026-04-22T04:37:45.261083Z","shell.execute_reply.started":"2026-04-22T04:37:12.394502Z","shell.execute_reply":"2026-04-22T04:37:45.260325Z"}},"outputs":[],"execution_count":null},{"cell_type":"code","source":"class EmbeddingBuilder:\n    \"\"\"\n    Nhận articles + transactions, sinh ra:\n      - df_items  với cột: item_embed_vec (list 160d), item_cluster_id\n      - df_users  với cột: user_embed_vec (list 160d), customer_cluster_id\n                           user_vec_7d, user_vec_14d, user_vec_30d (list 160d)\n \n    Sau đó lưu ra 2 file parquet để dùng lại.\n    \"\"\"\n \n    def __init__(\n        self,\n        articles: pd.DataFrame,\n        transactions: pd.DataFrame,  # hist_train hoặc hist_val\n        customers: pd.DataFrame,\n        image_dir: str = None,\n        n_item_clusters: int = 50,\n        n_user_clusters: int = 30,\n        device: str = \"cuda\",\n        out_dir: str = \"/kaggle/working/\",\n    ):\n        self.articles    = articles.copy()\n        self.trans       = transactions.copy()\n        self.customers   = customers.copy()\n        self.image_dir   = image_dir\n        self.n_item_clus = n_item_clusters\n        self.n_user_clus = n_user_clusters\n        self.device      = device\n        self.out_dir     = out_dir\n\n        # ── Fix article_id format ──────────────────────\n        self.articles[\"article_id\"] = self.articles[\"article_id\"].astype(str).str.zfill(10)\n        self.trans[\"article_id\"]    = self.trans[\"article_id\"].astype(str).str.zfill(10)\n        # ───────────────────────────────────────────────\n\n        self.item_embed  = None   # np.ndarray (n_items, 160)\n        self.art_id_list = articles[\"article_id\"].tolist()\n        self.art_idx     = {aid: i for i, aid in enumerate(self.art_id_list)}\n\n \n    # ── A2. Text vector ──────────────────────────────────────\n    def _build_text_vec(self):\n        def build_text(row):\n            parts = [\n                f\"Product: {row['prod_name']}\",\n                f\"Type: {row['product_type_name']}\",\n                f\"Group: {row['product_group_name']}\",\n                f\"Appearance: {row['graphical_appearance_name']}\",\n                f\"Colour: {row['perceived_colour_master_name']} {row['colour_group_name']}\",\n                f\"Department: {row['department_name']}\",\n                f\"Section: {row['section_name']}\",\n                f\"Garment: {row['garment_group_name']}\",\n                f\"Index: {row['index_group_name']} {row['index_name']}\",\n                f\"Description: {row['detail_desc']}\",\n            ]\n            return \" | \".join(p for p in parts if pd.notna(p.split(\": \", 1)[1]))\n    \n        texts = self.articles.apply(build_text, axis=1).tolist()\n    \n        model = SentenceTransformer(\"paraphrase-multilingual-MiniLM-L12-v2\")\n        raw   = model.encode(texts, batch_size=256, show_progress_bar=True,\n                             device=self.device)\n        svd   = TruncatedSVD(n_components=96, random_state=42)\n        return svd.fit_transform(raw)   # (n_items, 96)\n \n    # ── A3. Image vector (optional) ──────────────────────────\n    def _build_image_vec(self):\n        if self.image_dir is None:\n            print(\"  [skip] image_dir không được cung cấp → dùng zero vector\")\n            return np.zeros((len(self.art_id_list), 64))\n            \n        model_clip, preprocess = load(\"ViT-B/32\", device=self.device)\n \n        class _DS(Dataset):\n            def __init__(self, ids, img_dir, prep):\n                self.ids  = ids\n                self.dir  = img_dir\n                self.prep = prep\n            def __len__(self): return len(self.ids)\n            def __getitem__(self, idx):\n                aid    = str(int(self.ids[idx])).zfill(10)\n                folder = aid[:3]\n                path   = os.path.join(self.dir, f\"0{folder}\", f\"0{aid}.jpg\")\n                try:\n                    img = self.prep(PILImage.open(path).convert(\"RGB\"))\n                except Exception:\n                    img = torch.zeros(3, 224, 224)\n                return img, idx\n \n        ds     = _DS(self.art_id_list, self.image_dir, preprocess)\n        loader = DataLoader(ds, batch_size=256, num_workers=2, pin_memory=True)\n        raw    = np.zeros((len(self.art_id_list), 512))\n \n        model_clip.eval()\n        with torch.no_grad():\n            for imgs, idxs in loader:\n                feats = model_clip.encode_image(imgs.to(self.device))\n                raw[idxs] = feats.cpu().float().numpy()\n \n        svd = TruncatedSVD(n_components=64, random_state=42)\n        return svd.fit_transform(raw)   # (n_items, 64)\n \n    # ── A4. Build item embedding ─────────────────────────────\n    def build_item_embeddings(self):\n        print(\"Building item embeddings...\")\n        text_vec  = self._build_text_vec()   # (n, 96)\n        image_vec = self._build_image_vec()  # (n, 64)\n    \n        combined        = np.hstack([text_vec, image_vec])  # (n, 160)\n        self.item_embed = normalize(combined)\n    \n        km = KMeans(n_clusters=self.n_item_clus, random_state=42, n_init=10)\n        item_cluster = km.fit_predict(self.item_embed)\n        self.item_km_centers = km.cluster_centers_\n    \n        df_items = self.articles[[\"article_id\"]].copy()\n        df_items[\"item_embed_vec\"]  = list(self.item_embed)\n        df_items[\"item_cluster_id\"] = item_cluster.astype(\"int16\")\n    \n        out = self.out_dir + \"df_items.parquet\"\n        df_items.to_parquet(out, index=False)\n        print(f\"  Saved {out}  shape={df_items.shape}\")\n        return df_items\n     \n    # ── A5. Build user embedding ─────────────────────────────\n    def build_user_embeddings(self, df_items: pd.DataFrame):\n        print(\"Building user embeddings...\")\n        tr = self.trans.copy()\n        tr[\"t_dat\"] = pd.to_datetime(tr[\"t_dat\"])\n        max_date = tr[\"t_dat\"].max()\n        tr[\"days_since\"] = (max_date - tr[\"t_dat\"]).dt.days + 1\n    \n        # Map article_id → row index trong item_embed matrix\n        DIM = 160\n        item_vec_matrix = np.vstack(df_items[\"item_embed_vec\"].values)  # (n_items, 160)\n        art_idx = {aid: i for i, aid in enumerate(df_items[\"article_id\"])}\n    \n        # Chỉ giữ transaction có article tồn tại\n        tr[\"item_idx\"] = tr[\"article_id\"].map(art_idx)\n\n        # ── DEBUG ──────────────────────────────────────\n        print(f\"tr article_id dtype    : {tr['article_id'].dtype}\")\n        print(f\"df_items article_id dtype: {df_items['article_id'].dtype}\")\n        print(f\"tr article_id sample   : {tr['article_id'].iloc[0]!r}\")\n        print(f\"df_items article_id sample: {df_items['article_id'].iloc[0]!r}\")\n        print(f\"tr shape trước dropna  : {tr.shape}\")\n        print(f\"item_idx null count    : {tr['item_idx'].isna().sum()}\")\n        # ───────────────────────────────────────────────\n        \n        tr = tr.dropna(subset=[\"item_idx\"])\n        tr[\"item_idx\"] = tr[\"item_idx\"].astype(int)\n    \n        # Map customer_id → integer index\n        customers = tr[\"customer_id\"].unique()\n        cust_idx  = {c: i for i, c in enumerate(customers)}\n        tr[\"cust_idx\"] = tr[\"customer_id\"].map(cust_idx)\n        n_users = len(customers)\n    \n        def _vectorized_weighted_mean(mask=None):\n            \"\"\"\n            Tính weighted mean cho tất cả users cùng lúc dùng sparse matrix.\n            weight = 1 / days_since, normalize theo user.\n            \"\"\"\n            sub = tr if mask is None else tr[mask]\n            if len(sub) == 0:\n                return np.zeros((n_users, DIM))\n    \n            weights = 1.0 / sub[\"days_since\"].values  # (n_rows,)\n    \n            # Sparse matrix: (n_users, n_items), value = sum of weights\n            W = csr_matrix(\n                (weights, (sub[\"cust_idx\"].values, sub[\"item_idx\"].values)),\n                shape=(n_users, len(art_idx))\n            )\n    \n            # Normalize mỗi hàng (user) về tổng = 1\n            row_sums = np.asarray(W.sum(axis=1)).ravel()  # (n_users,)\n            row_sums[row_sums == 0] = 1  # tránh chia 0\n            W = W.multiply(1.0 / row_sums[:, None])\n    \n            # (n_users, n_items) @ (n_items, 160) = (n_users, 160)\n            return W.dot(item_vec_matrix)\n    \n        print(\"  Computing user_embed_vec ...\")\n        user_vecs     = _vectorized_weighted_mean()\n        print(\"  Computing user_vec_7d ...\")\n        user_vecs_7d  = _vectorized_weighted_mean(tr[\"days_since\"] <= 7)\n        print(\"  Computing user_vec_14d ...\")\n        user_vecs_14d = _vectorized_weighted_mean(tr[\"days_since\"] <= 14)\n        print(\"  Computing user_vec_30d ...\")\n        user_vecs_30d = _vectorized_weighted_mean(tr[\"days_since\"] <= 30)\n    \n        df_users = pd.DataFrame({\n            \"customer_id\":    customers,\n            \"user_embed_vec\": list(user_vecs),\n            \"user_vec_7d\":    list(user_vecs_7d),\n            \"user_vec_14d\":   list(user_vecs_14d),\n            \"user_vec_30d\":   list(user_vecs_30d),\n        })\n    \n        # Merge demographic features (giữ nguyên như cũ)\n        demo_cols = [\"customer_id\", \"FN\", \"Active\", \"club_member_status\",\n                     \"fashion_news_frequency\", \"age\"]\n        demo = self.customers[demo_cols].copy()\n        demo[\"FN\"]     = demo[\"FN\"].fillna(0).astype(int)\n        demo[\"Active\"] = demo[\"Active\"].fillna(0).astype(int)\n        demo[\"club_member_status\"]    = demo[\"club_member_status\"].map(\n            {\"ACTIVE\": 2, \"PRE-CREATE\": 1, \"LEFT CLUB\": -1}).fillna(0).astype(int)\n        demo[\"fashion_news_frequency\"] = demo[\"fashion_news_frequency\"].map(\n            {\"Regularly\": 2, \"Monthly\": 1, \"NONE\": 0, \"None\": 0}).fillna(0).astype(int)\n        demo[\"age\"] = demo[\"age\"].fillna(demo[\"age\"].median()).astype(\"int8\")\n        bins = [0, 25, 35, 45, 55, 65, 999]\n        demo[\"age_binning\"] = pd.cut(demo[\"age\"], bins=bins,\n                                     labels=range(6), right=True).astype(\"int8\")\n        df_users = df_users.merge(demo, on=\"customer_id\", how=\"left\")\n    \n        # KMeans cluster\n        vecs = np.vstack(df_users[\"user_embed_vec\"].values)\n        km   = KMeans(n_clusters=self.n_user_clus, random_state=42, n_init=10)\n        df_users[\"customer_cluster_id\"] = km.fit_predict(vecs).astype(\"int16\")\n        self.user_km_centers = km.cluster_centers_\n    \n        out = self.out_dir + \"df_users.parquet\"\n        df_users.to_parquet(out, index=False)\n        print(f\"  Saved {out}  shape={df_users.shape}\")\n        return df_users\n \n    # ── Entry point ──────────────────────────────────────────\n    def run(self):\n        df_items = self.build_item_embeddings()\n        df_users = self.build_user_embeddings(df_items)\n        return df_items, df_users","metadata":{"trusted":true,"execution":{"iopub.status.busy":"2026-04-22T04:37:45.262090Z","iopub.execute_input":"2026-04-22T04:37:45.262864Z","iopub.status.idle":"2026-04-22T04:37:45.289639Z","shell.execute_reply.started":"2026-04-22T04:37:45.262824Z","shell.execute_reply":"2026-04-22T04:37:45.288774Z"}},"outputs":[],"execution_count":null},{"cell_type":"code","source":"class CopurchaseFeatureBuilder:\n    def __init__(self, df_items, df_users, transactions, articles):\n        self.df_items = df_items.set_index(\"article_id\")\n        self.df_users = df_users.set_index(\"customer_id\")\n        self.articles = articles.set_index(\"article_id\")\n\n        tr = transactions.copy()\n        tr[\"t_dat\"] = pd.to_datetime(tr[\"t_dat\"])\n        self.trans    = tr\n        self.max_date = tr[\"t_dat\"].max()\n\n        self._seq_prob           = None\n        self._group_seq_prob     = None\n        self._cluster_copurchase = None\n        self._item_km_centers    = None\n\n    # ── B1. Precompute sequence probability ──────────────────\n    def precompute_sequence(self):\n        print(\"Precomputing sequence probabilities...\")\n        tr_sorted = self.trans.sort_values([\"customer_id\", \"t_dat\"])\n        tr_sorted[\"next_article\"] = tr_sorted.groupby(\"customer_id\")[\"article_id\"].shift(-1)\n        tr_sorted = tr_sorted.dropna(subset=[\"next_article\"])\n\n        counts = tr_sorted.groupby([\"article_id\", \"next_article\"]).size().reset_index(name=\"cnt\")\n        totals = counts.groupby(\"article_id\")[\"cnt\"].sum().reset_index(name=\"total\")\n        counts = counts.merge(totals, on=\"article_id\")\n        counts[\"prob\"] = counts[\"cnt\"] / counts[\"total\"]\n        self._seq_prob = counts.set_index([\"article_id\", \"next_article\"])[\"prob\"].to_dict()\n\n        art_group = self.articles[\"product_group_name\"].to_dict()\n        tr_sorted[\"group\"]      = tr_sorted[\"article_id\"].map(art_group)\n        tr_sorted[\"next_group\"] = tr_sorted[\"next_article\"].map(art_group)\n        tr_sorted = tr_sorted.dropna(subset=[\"group\", \"next_group\"])\n\n        gcounts = tr_sorted.groupby([\"group\", \"next_group\"]).size().reset_index(name=\"cnt\")\n        gtotals = gcounts.groupby(\"group\")[\"cnt\"].sum().reset_index(name=\"total\")\n        gcounts = gcounts.merge(gtotals, on=\"group\")\n        gcounts[\"prob\"] = gcounts[\"cnt\"] / gcounts[\"total\"]\n        self._group_seq_prob = gcounts.set_index([\"group\", \"next_group\"])[\"prob\"].to_dict()\n\n        gc.collect()\n        print(\"  Done.\")\n\n    # ── B2. Precompute cluster co-purchase (vectorized) ──────\n    def precompute_cluster_copurchase(self):\n        print(\"Precomputing cluster copurchase...\")\n        from scipy.sparse import csr_matrix\n\n        cust_cluster = self.df_users[\"customer_cluster_id\"].to_dict()\n        self.trans[\"user_cluster\"] = self.trans[\"customer_id\"].map(cust_cluster)\n        tr = self.trans.dropna(subset=[\"user_cluster\"]).copy()\n        tr[\"user_cluster\"] = tr[\"user_cluster\"].astype(int)\n\n        art_ids = self.df_items.index.tolist()\n        art_idx = {aid: i for i, aid in enumerate(art_ids)}\n        item_vec_matrix = np.vstack(self.df_items[\"item_embed_vec\"].values)\n\n        tr[\"item_idx\"] = tr[\"article_id\"].map(art_idx)\n        tr = tr.dropna(subset=[\"item_idx\"])\n        tr[\"item_idx\"] = tr[\"item_idx\"].astype(int)\n\n        n_clusters = int(tr[\"user_cluster\"].max()) + 1\n        ones = np.ones(len(tr))\n        W = csr_matrix(\n            (ones, (tr[\"user_cluster\"].values, tr[\"item_idx\"].values)),\n            shape=(n_clusters, len(art_ids))\n        )\n        row_sums = np.asarray(W.sum(axis=1)).ravel()\n        row_sums[row_sums == 0] = 1\n        W = W.multiply(1.0 / row_sums[:, None])\n        self._cluster_copurchase = W.dot(item_vec_matrix).astype(\"float32\")\n\n        gc.collect()\n        print(\"  Done.\")\n\n    # ── B3. Cluster popularity (7/14/30d) ────────────────────\n    def _build_cluster_pop(self, days):\n        cutoff = self.max_date - pd.Timedelta(days=days)\n        tr_win = self.trans[self.trans[\"t_dat\"] > cutoff].copy()\n        cust_cluster = self.df_users[\"customer_cluster_id\"].to_dict()\n        tr_win[\"user_cluster\"] = tr_win[\"customer_id\"].map(cust_cluster)\n\n        pop = (\n            tr_win.dropna(subset=[\"user_cluster\"])\n            .groupby([\"user_cluster\", \"article_id\"])\n            .size()\n            .reset_index(name=\"cnt\")\n        )\n        totals = pop.groupby(\"user_cluster\")[\"cnt\"].sum().reset_index(name=\"total\")\n        pop    = pop.merge(totals, on=\"user_cluster\")\n        pop[\"score\"] = pop[\"cnt\"] / pop[\"total\"]\n        return pop.set_index([\"user_cluster\", \"article_id\"])[\"score\"].to_dict()\n\n    # ── B4. Add features (batch + GPU) ───────────────────────\n    def add_features(self, candidates: pd.DataFrame, batch_size: int = 50000) -> pd.DataFrame:\n        import torch\n        dev = torch.device(\"cuda\" if torch.cuda.is_available() else \"cpu\")\n\n        if self._seq_prob is None:\n            self.precompute_sequence()\n        if self._cluster_copurchase is None:\n            self.precompute_cluster_copurchase()\n\n        cp7  = self._build_cluster_pop(7)\n        cp14 = self._build_cluster_pop(14)\n        cp30 = self._build_cluster_pop(30)\n\n        last_bought = (\n            self.trans.sort_values(\"t_dat\")\n            .groupby(\"customer_id\")[\"article_id\"].last().to_dict()\n        )\n        art_group = self.articles[\"product_group_name\"].to_dict()\n        DIM = 160\n\n        # Chuẩn bị numpy arrays trên CPU\n        item_embed_np  = np.vstack(self.df_items[\"item_embed_vec\"].values).astype(\"float32\")\n        user_embed_np  = np.vstack(self.df_users[\"user_embed_vec\"].values).astype(\"float32\")\n        user_7d_np     = np.vstack(self.df_users[\"user_vec_7d\"].values).astype(\"float32\")\n        user_14d_np    = np.vstack(self.df_users[\"user_vec_14d\"].values).astype(\"float32\")\n        user_30d_np    = np.vstack(self.df_users[\"user_vec_30d\"].values).astype(\"float32\")\n        cluster_cop_np = self._cluster_copurchase  # (n_clusters, 160)\n        item_km_np     = self._item_km_centers.astype(\"float32\") \\\n                         if self._item_km_centers is not None else None\n        item_clus_arr  = self.df_items[\"item_cluster_id\"].values.astype(\"int32\")\n        user_clus_arr  = self.df_users[\"customer_cluster_id\"].values.astype(\"int32\")\n\n        item_id_to_idx = {aid: i for i, aid in enumerate(self.df_items.index)}\n        user_id_to_idx = {uid: i for i, uid in enumerate(self.df_users.index)}\n\n        def batch_cosine_gpu(A_np, B_np):\n            A = torch.from_numpy(A_np).to(dev)\n            B = torch.from_numpy(B_np).to(dev)\n            res = ((A / A.norm(dim=1, keepdim=True).clamp(1e-9)) *\n                   (B / B.norm(dim=1, keepdim=True).clamp(1e-9))).sum(dim=1)\n            out = res.cpu().numpy().astype(\"float32\")\n            del A, B, res\n            torch.cuda.empty_cache()\n            return out\n\n        results = []\n        n = len(candidates)\n        print(f\"Computing features for {n:,} rows (batch={batch_size}, device={dev})...\")\n\n        for start in range(0, n, batch_size):\n            end  = min(start + batch_size, n)\n            cand = candidates.iloc[start:end].copy()\n            print(f\"  Batch {start:,}:{end:,}...\")\n\n            cids = cand[\"customer_id\"].values\n            aids = cand[\"article_id\"].values\n\n            i_idx  = np.array([item_id_to_idx.get(a, 0) for a in aids])\n            u_idx  = np.array([user_id_to_idx.get(c, 0) for c in cids])\n            u_clus = user_clus_arr[u_idx].clip(0)\n\n            I   = item_embed_np[i_idx]\n            U   = user_embed_np[u_idx]\n            U7  = user_7d_np[u_idx]\n            U14 = user_14d_np[u_idx]\n            U30 = user_30d_np[u_idx]\n\n            # cosine_sim_to_history / recent_7d / 14d / 30d\n            cand[\"cosine_sim_to_history\"]    = batch_cosine_gpu(U,   I)\n            cand[\"cosine_sim_to_recent_7d\"]  = batch_cosine_gpu(U7,  I)\n            cand[\"cosine_sim_to_recent_14d\"] = batch_cosine_gpu(U14, I)\n            cand[\"cosine_sim_to_recent_30d\"] = batch_cosine_gpu(U30, I)\n\n            # nn_copurchase_score\n            last_arts = [last_bought.get(c) for c in cids]\n            la_idx    = np.array([item_id_to_idx.get(a, 0) if a is not None else 0\n                                  for a in last_arts])\n            L = item_embed_np[la_idx]\n            L[np.array([a is None for a in last_arts])] = 0\n            cand[\"nn_copurchase_score\"] = batch_cosine_gpu(L, I)\n            del L\n\n            # is_next_in_sequence\n            cand[\"is_next_in_sequence\"] = np.array([\n                self._seq_prob.get((la, a), 0.0) if la is not None else 0.0\n                for la, a in zip(last_arts, aids)\n            ], dtype=\"float32\")\n\n            # is_group_next_in_sequence\n            last_groups = [art_group.get(la) for la in last_arts]\n            item_groups = [art_group.get(a)  for a  in aids]\n            cand[\"is_group_next_in_sequence\"] = np.array([\n                self._group_seq_prob.get((lg, ig), 0.0) if lg is not None else 0.0\n                for lg, ig in zip(last_groups, item_groups)\n            ], dtype=\"float32\")\n\n            # style_cluster_score\n            if item_km_np is not None:\n                ic = item_clus_arr[i_idx].clip(0)\n                C  = item_km_np[ic]\n                cand[\"style_cluster_score\"] = batch_cosine_gpu(U, C)\n                del C\n            else:\n                cand[\"style_cluster_score\"] = np.float32(0.0)\n\n            # cluster_copurchase_score\n            CM = cluster_cop_np[u_clus]\n            cand[\"cluster_copurchase_score\"] = batch_cosine_gpu(CM, I)\n            del CM\n\n            # cluster_pop_7d / 14d / 30d\n            cand[\"cluster_pop_7d\"]  = np.array([cp7.get( (c, a), 0.0) for c, a in zip(u_clus, aids)], dtype=\"float32\")\n            cand[\"cluster_pop_14d\"] = np.array([cp14.get((c, a), 0.0) for c, a in zip(u_clus, aids)], dtype=\"float32\")\n            cand[\"cluster_pop_30d\"] = np.array([cp30.get((c, a), 0.0) for c, a in zip(u_clus, aids)], dtype=\"float32\")\n\n            results.append(cand)\n            del I, U, U7, U14, U30\n            gc.collect()\n\n        out = pd.concat(results, ignore_index=True)\n        print(f\"  Done. Shape: {out.shape}\")\n        return out\n\n    def set_item_km_centers(self, centers):\n        self._item_km_centers = centers","metadata":{"trusted":true,"execution":{"iopub.status.busy":"2026-04-22T04:37:45.290910Z","iopub.execute_input":"2026-04-22T04:37:45.291299Z","iopub.status.idle":"2026-04-22T04:37:45.352815Z","shell.execute_reply.started":"2026-04-22T04:37:45.291271Z","shell.execute_reply":"2026-04-22T04:37:45.352179Z"}},"outputs":[],"execution_count":null},{"cell_type":"code","source":"def process_and_save_by_chunks(input_path, output_path, builder, chunk_size=10_000_000):\n    \"\"\"\n    Đọc file 119M dòng theo từng chunk 10M, xử lý và ghi xuống đĩa ngay.\n    \"\"\"\n    # Sử dụng Polars để scan file (không tốn RAM)\n    lazy_df = pl.scan_parquet(input_path)\n    total_rows = 119121000 # Số dòng bạn đã check\n    \n    first_batch = True\n    \n    for start in range(0, total_rows, chunk_size):\n        print(f\"--- Processing chunk: {start} to {start + chunk_size} ---\")\n        \n        # Load 1 chunk vào RAM (Pandas)\n        batch_df = lazy_df.slice(start, chunk_size).collect().to_pandas()\n        \n        # Tiền xử lý ID\n        batch_df[\"article_id\"] = batch_df[\"article_id\"].astype(str).str.zfill(10)\n        batch_df[\"customer_id\"] = batch_df[\"customer_id\"].astype(str)\n        \n        # Chạy Feature Engineering\n        featured_batch = builder.add_features(batch_df, batch_size=100000)\n        \n        # Ghi xuống đĩa (Dùng fastparquet hoặc pyarrow để append)\n        if first_batch:\n            featured_batch.to_parquet(output_path, engine='pyarrow', index=False)\n            first_batch = False\n        else:\n            # Append vào file đã có\n            table = pa.Table.from_pandas(featured_batch)\n            with pq.ParquetWriter(output_path, table.schema, engine='pyarrow', append=True) as writer:\n                writer.write_table(table)\n        \n        del batch_df, featured_batch\n        gc.collect()","metadata":{"trusted":true,"execution":{"iopub.status.busy":"2026-04-22T04:37:45.353910Z","iopub.execute_input":"2026-04-22T04:37:45.354557Z","iopub.status.idle":"2026-04-22T04:37:45.368810Z","shell.execute_reply.started":"2026-04-22T04:37:45.354527Z","shell.execute_reply":"2026-04-22T04:37:45.367959Z"}},"outputs":[],"execution_count":null},{"cell_type":"code","source":"# --- CẤU HÌNH ĐƯỜNG DẪN ---\nBASE  = \"/kaggle/input/datasets/huyenpham22/h-and-m-0-4-stratified-sample\"\nEMBED = \"/kaggle/input/datasets/huyenpham22/h-and-m-val-train-embed\"\nWORK  = \"/kaggle/working/\"\nINPUT_CANDIDATES = BASE + \"/test_candidates_top300.parquet\"\nOUTPUT_PATH = WORK + \"test_candidates_with_features.parquet\"\n\n# --- 1. LOAD CÁC TẬP DỮ LIỆU BỔ TRỢ (NHỎ HƠN) ---\nprint(\"Loading supporting datasets...\")\n\n# Với Parquet, load xong rồi mới ép kiểu (zfill)\narticles  = pd.read_csv(BASE + \"/articles.csv\", dtype={\"article_id\": \"string\"})\narticles[\"article_id\"] = articles[\"article_id\"].str.zfill(10)\n\ncustomers = pd.read_csv(BASE + \"/customers.csv\", dtype={\"customer_id\": \"string\"})\n\n# SỬA LỖI TẠI ĐÂY: Load hist_train và ép kiểu thủ công\nhist_test  = pd.read_parquet(\"/kaggle/input/datasets/huyenpham22/h-and-m-0-4-stratified-sample/transactions_train.parquet\")\nhist_test[\"article_id\"] = hist_test[\"article_id\"].astype(str).str.zfill(10)\nhist_test[\"customer_id\"] = hist_test[\"customer_id\"].astype(str)\n\ndf_items       = pd.read_parquet(\"/kaggle/input/datasets/huyenpham22/h-and-m-val-train-embed/df_items.parquet\")\ndf_users_test        = pd.read_parquet(\"/kaggle/input/datasets/huyenpham22/h-and-m-val-train-embed/test_embed/test_embed/df_users.parquet\")\nitem_km_centers_train = np.load(\"/kaggle/input/datasets/huyenpham22/h-and-m-val-train-embed/test_embed/test_embed/item_km_centers.npy\")\n\n# --- 2. KHỞI TẠO BUILDER ---\nfeat_builder_train = CopurchaseFeatureBuilder(\n    df_items=df_items,\n    df_users=df_users_test,\n    transactions=hist_test,\n    articles=articles,\n)\nfeat_builder_train.set_item_km_centers(item_km_centers_train)\n\n# Thực hiện precompute một lần duy nhất trước vòng lặp\nfeat_builder_train.precompute_sequence()\nfeat_builder_train.precompute_cluster_copurchase()\n\n# --- 3. XỬ LÝ THEO CHUNK (GIẢI PHÁP CHO 119M DÒNG) ---\n# Dùng Polars scan để không tốn RAM load toàn bộ file\nmetadata = pq.read_metadata(INPUT_CANDIDATES)\ntotal_rows = metadata.num_rows\n\nlazy_df = pl.scan_parquet(INPUT_CANDIDATES)\nchunk_size = 10_000_000 # Xử lý 10 triệu dòng mỗi đợt\n\nprint(f\"Bắt đầu xử lý {total_rows} dòng theo từng đợt {chunk_size}...\")\n\nwriter = None\n\nfor start in range(0, total_rows, chunk_size):\n    # Lấy một đoạn dữ liệu và chuyển sang Pandas\n    batch_cand = lazy_df.slice(start, chunk_size).collect().to_pandas()\n    \n    if batch_cand.empty:\n        break\n        \n    # Chuẩn hóa định dạng ID để join/map chính xác\n    batch_cand[\"article_id\"] = batch_cand[\"article_id\"].astype(str).str.zfill(10)\n    batch_cand[\"customer_id\"] = batch_cand[\"customer_id\"].astype(str)\n    \n    # Tính toán features (sử dụng GPU bên trong builder)\n    batch_feat = feat_builder_train.add_features(batch_cand)\n    \n    # Ghi nối đuôi (Append) trực tiếp xuống ổ cứng bằng ParquetWriter\n    table = pa.Table.from_pandas(batch_feat)\n    if writer is None:\n        # Khởi tạo writer tại batch đầu tiên với schema của batch đó\n        writer = pq.ParquetWriter(OUTPUT_PATH, table.schema, compression='snappy')\n    \n    writer.write_table(table)\n    \n    print(f\"Tiến độ: {min(start + chunk_size, total_rows):,} / {total_rows:,} dòng\")\n    \n    # GIẢI PHÓNG BỘ NHỚ TRIỆT ĐỂ SAU MỖI BATCH\n    del batch_cand, batch_feat, table\n    gc.collect()\n    torch.cuda.empty_cache()\n\n# Đóng writer để hoàn tất file\nif writer:\n    writer.close()\n\nprint(f\"\\n[Xong] File kết quả train lưu tại: {OUTPUT_PATH}\")","metadata":{"trusted":true,"execution":{"iopub.status.busy":"2026-04-22T04:37:45.370077Z","iopub.execute_input":"2026-04-22T04:37:45.370447Z","iopub.status.idle":"2026-04-22T04:42:46.483458Z","shell.execute_reply.started":"2026-04-22T04:37:45.370421Z","shell.execute_reply":"2026-04-22T04:42:46.482400Z"}},"outputs":[],"execution_count":null},{"cell_type":"code","source":"","metadata":{"trusted":true},"outputs":[],"execution_count":null}]}