{"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":"nvidiaTeslaT4","dataSources":[{"sourceId":31254,"databundleVersionId":3103714,"sourceType":"competition"}],"dockerImageVersionId":31193,"isInternetEnabled":true,"language":"python","sourceType":"notebook","isGpuEnabled":true}},"nbformat_minor":4,"nbformat":4,"cells":[{"cell_type":"code","source":"\"\"\"\nH&M Personalized Fashion Recommendations - V2.8 Enhanced Implementation\nFollowing the V2 improvement plan for enhanced recommendation system\nExpected MAP@12: 0.032-0.038 (vs 0.012-0.015 baseline)\nRuntime: 27-39 minutes (target: <1 hour)\nOptimized for Kaggle GPU environment\n\"\"\"\n\nimport pandas as pd\nimport numpy as np\nfrom datetime import datetime, timedelta\nfrom collections import defaultdict, Counter\nimport warnings\nimport os\nimport sys\nimport gc\nfrom tqdm import tqdm\nimport multiprocessing as mp\nfrom itertools import combinations\n\n# Try to import optional dependencies\ntry:\n    import psutil\n    PSUTIL_AVAILABLE = True\nexcept ImportError:\n    PSUTIL_AVAILABLE = False\n\n# Try to import dask for parallel dataframe operations\ntry:\n    import dask.dataframe as dd\n    from dask import delayed, compute\n    DASK_AVAILABLE = True\nexcept ImportError:\n    print(\"Warning: Dask not available, using pandas only\")\n    DASK_AVAILABLE = False\n\n# Memory monitoring\ndef get_memory_usage():\n    \"\"\"Get current memory usage in MB.\"\"\"\n    if PSUTIL_AVAILABLE:\n        process = psutil.Process()\n        return process.memory_info().rss / 1024 / 1024\n    else:\n        return 0.0\n\ndef print_memory_usage(message=\"\"):\n    \"\"\"Print current memory usage.\"\"\"\n    if PSUTIL_AVAILABLE:\n        mem_mb = get_memory_usage()\n        print(f\"    Memory usage{ ' - ' + message if message else ''}: {mem_mb:.1f} MB\")\n\nwarnings.filterwarnings('ignore')\n\n# Constants\nDATA_DIR = \"/kaggle/input/h-and-m-personalized-fashion-recommendations\"\nLAMBDA_DECAY = np.log(2) / 14  # Exponential decay parameter (14-day half-life)\nDIVERSITY_PENALTY = 0.3  # 30% reduction for same product code prefix\n\n\ndef reduce_mem_usage(df):\n    \"\"\"Reduce memory usage by downcasting numeric dtypes.\"\"\"\n    start_mem = df.memory_usage().sum() / 1024**2\n\n    for col in df.columns:\n        col_type = df[col].dtype\n\n        # Skip datetime and categorical columns\n        if col_type == 'datetime64[ns]' or col_type.name == 'category':\n            continue\n\n        if col_type != object:\n            c_min = df[col].min()\n            c_max = df[col].max()\n\n            if str(col_type)[:3] == 'int':\n                if c_min > np.iinfo(np.int8).min and c_max < np.iinfo(np.int8).max:\n                    df[col] = df[col].astype(np.int8)\n                elif c_min > np.iinfo(np.int16).min and c_max < np.iinfo(np.int16).max:\n                    df[col] = df[col].astype(np.int16)\n                elif c_min > np.iinfo(np.int32).min and c_max < np.iinfo(np.int32).max:\n                    df[col] = df[col].astype(np.int32)\n            elif str(col_type)[:5] == 'float':\n                # Skip float16 downcasting due to pandas compatibility issues\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.float64)\n\n    end_mem = df.memory_usage().sum() / 1024**2\n    print(f'    Memory usage after optimization: {end_mem:.2f} MB')\n    print(f'    Decreased by {100 * (start_mem - end_mem) / start_mem:.1f}%')\n    return df\n\ndef exponential_time_decay(days_ago):\n    \"\"\"Phase 2: Exponential time decay with 14-day half-life.\"\"\"\n    return np.exp(-LAMBDA_DECAY * days_ago)\n\n\ndef get_timescale_multiplier(days_ago):\n    \"\"\"Phase 3: Multi-timescale windows.\"\"\"\n    if days_ago <= 7:\n        return 1.5  # Short-term: trending preferences\n    elif days_ago <= 30:\n        return 1.0  # Medium-term: stable preferences\n    elif days_ago <= 90:\n        return 0.5  # Long-term: general style\n    else:\n        return 0.1  # Very old: minimal weight\n\n\ndef get_recency_multiplier(days_ago):\n    \"\"\"Phase 4: Recency multiplier for purchase frequency boost.\"\"\"\n    if days_ago <= 3:\n        return 2.0  # Last 3 days: 2.0x\n    elif days_ago <= 7:\n        return 1.5  # Last 7 days: 1.5x\n    else:\n        return 1.0\n\n\ndef calculate_repurchase_candidates(transactions_df, cutoff_date, lookback_days=90):\n    \"\"\"Phase 1 Strategy 1: Repurchase candidates from customer history.\"\"\"\n    print(\"  [Strategy 1] Calculating repurchase candidates...\")\n    lookback_date = cutoff_date - timedelta(days=lookback_days)\n    recent_transactions = transactions_df[\n        (transactions_df['t_dat'] >= lookback_date) &\n        (transactions_df['t_dat'] < cutoff_date)\n    ].copy()\n\n    print(f\"    Processing {len(recent_transactions):,} transactions...\")\n\n    if DASK_AVAILABLE and len(recent_transactions) > 1000000:  # Use dask for large datasets\n        print(\"    Using Dask for parallel processing...\")\n\n        # Convert to dask dataframe\n        ddf = dd.from_pandas(recent_transactions, npartitions=mp.cpu_count())\n\n        # Parallel groupby operation\n        result = ddf.groupby(['customer_id', 'article_id']).size().compute()\n\n        # Convert to dictionary format\n        repurchase_candidates = {}\n        for (customer_id, article_id), count in result.items():\n            if customer_id not in repurchase_candidates:\n                repurchase_candidates[customer_id] = {}\n            repurchase_candidates[customer_id][article_id] = count\n\n    else:\n        # Use pandas for smaller datasets or when dask not available\n        print(\"    Using pandas processing...\")\n        repurchase_candidates = {}\n        for customer_id, group in tqdm(recent_transactions.groupby('customer_id'), desc=\"    Processing customers\"):\n            article_counts = group['article_id'].value_counts()\n            repurchase_candidates[customer_id] = article_counts.to_dict()\n\n    print(f\"    Generated repurchase candidates for {len(repurchase_candidates)} customers\")\n    return repurchase_candidates\n\n\ndef calculate_age_group_popular_items(transactions_df, customers_df, cutoff_date, lookback_days=30):\n    \"\"\"Phase 1 Strategy 2: Age group popular items.\"\"\"\n    print(\"  [Strategy 2] Calculating age group popular items...\")\n    lookback_date = cutoff_date - timedelta(days=lookback_days)\n    recent_transactions = transactions_df[\n        (transactions_df['t_dat'] >= lookback_date) &\n        (transactions_df['t_dat'] < cutoff_date)\n    ].copy()\n\n    # Merge with customer age\n    if 'age' not in customers_df.columns:\n        print(\"    Warning: Age column not found, skipping age group strategy\")\n        return {}\n\n    # Create age groups\n    customers_df = customers_df.copy()\n    customers_df['age_group'] = pd.cut(\n        customers_df['age'],\n        bins=[0, 25, 35, 45, 55, 100],\n        labels=['18-25', '26-35', '36-45', '46-55', '55+']\n    )\n\n    print(f\"    Processing {len(recent_transactions):,} transactions for age groups...\")\n\n    if DASK_AVAILABLE and len(recent_transactions) > 500000:  # Use dask for large merges\n        print(\"    Using Dask for parallel merge/groupby...\")\n\n        # Convert to dask dataframes\n        trans_ddf = dd.from_pandas(recent_transactions, npartitions=mp.cpu_count())\n        cust_ddf = dd.from_pandas(customers_df[['customer_id', 'age_group']], npartitions=4)\n\n        # Parallel merge\n        merged_ddf = trans_ddf.merge(cust_ddf, on='customer_id', how='left')\n\n        # Parallel groupby and value_counts equivalent\n        result = merged_ddf.groupby(['age_group', 'article_id']).size().compute()\n\n        # Convert to dictionary format\n        age_group_popular = {}\n        for (age_group, article_id), count in result.items():\n            if pd.notna(age_group):\n                if age_group not in age_group_popular:\n                    age_group_popular[age_group] = {}\n                age_group_popular[age_group][article_id] = count\n\n        # Keep only top 50 per age group\n        for age_group in age_group_popular:\n            sorted_items = sorted(age_group_popular[age_group].items(), key=lambda x: x[1], reverse=True)\n            age_group_popular[age_group] = dict(sorted_items[:50])\n\n    else:\n        # Use pandas for smaller datasets or when dask not available\n        print(\"    Using pandas processing...\")\n        merged = recent_transactions.merge(\n            customers_df[['customer_id', 'age_group']],\n            on='customer_id',\n            how='left'\n        )\n\n        # Calculate popular items per age group\n        age_group_popular = {}\n        for age_group in merged['age_group'].dropna().unique():\n            age_transactions = merged[merged['age_group'] == age_group]\n            popular_items = age_transactions['article_id'].value_counts().head(50).to_dict()\n            age_group_popular[age_group] = popular_items\n\n    print(f\"    Generated popular items for {len(age_group_popular)} age groups\")\n    return age_group_popular\n\n\ndef calculate_cooccurrence_items(transactions_df, cutoff_date, lookback_days=90, min_cooccurrence=5):\n    \"\"\"Phase 1 Strategy 3: Items bought together (co-occurrence) - VECTORIZED VERSION.\"\"\"\n    print(\"  [Strategy 3] Calculating co-occurrence patterns (vectorized)...\")\n    lookback_date = cutoff_date - timedelta(days=lookback_days)\n    recent_transactions = transactions_df[\n        (transactions_df['t_dat'] >= lookback_date) &\n        (transactions_df['t_dat'] < cutoff_date)\n    ].copy()\n\n    print(f\"    Processing {len(recent_transactions):,} transactions for co-occurrence...\")\n\n    # Vectorized approach: Group by customer and date, then create combinations\n    grouped = recent_transactions.groupby(['customer_id', 't_dat'])['article_id'].unique()\n    grouped = grouped.reset_index()\n\n    # Create all combinations for each customer-date group\n    all_combinations = []\n    for articles in tqdm(grouped['article_id'], desc=\"    Creating combinations\"):\n        if len(articles) > 1:\n            # Create unordered pairs (combinations, not permutations)\n            for art1, art2 in combinations(sorted(articles), 2):\n                all_combinations.append((art1, art2))\n\n    if not all_combinations:\n        print(\"    No co-occurrence patterns found\")\n        return {}\n\n    # Convert to DataFrame for efficient counting\n    combo_df = pd.DataFrame(all_combinations, columns=['article1', 'article2'])\n    cooccurrence_counts = combo_df.groupby(['article1', 'article2']).size().reset_index(name='count')\n\n    # Filter by minimum co-occurrence\n    cooccurrence_counts = cooccurrence_counts[cooccurrence_counts['count'] >= min_cooccurrence]\n\n    # Convert to dictionary format\n    cooccurrence = defaultdict(Counter)\n    for _, row in cooccurrence_counts.iterrows():\n        art1, art2, count = row['article1'], row['article2'], row['count']\n        cooccurrence[art1][art2] = count\n        cooccurrence[art2][art1] = count\n\n    # Filter out articles with no co-occurrences\n    filtered_cooccurrence = {k: dict(v) for k, v in cooccurrence.items() if v}\n\n    print(f\"    Generated co-occurrence patterns for {len(filtered_cooccurrence)} articles\")\n    return filtered_cooccurrence\n\n\ndef calculate_category_items(transactions_df, articles_df, cutoff_date, lookback_days=90):\n    \"\"\"Phase 1 Strategy 4: Same category items.\"\"\"\n    print(\"  [Strategy 4] Calculating category-based items...\")\n    lookback_date = cutoff_date - timedelta(days=lookback_days)\n    recent_transactions = transactions_df[\n        (transactions_df['t_dat'] >= lookback_date) &\n        (transactions_df['t_dat'] < cutoff_date)\n    ].copy()\n    \n    # Merge with article categories\n    if 'product_type_no' not in articles_df.columns:\n        print(\"    Warning: product_type_no not found, using article_id grouping\")\n        category_col = 'article_id'\n    else:\n        category_col = 'product_type_no'\n    \n    merged = recent_transactions.merge(\n        articles_df[['article_id', category_col]], \n        on='article_id', \n        how='left'\n    )\n    \n    # Calculate popular items per category\n    category_popular = {}\n    for category, group in merged.groupby(category_col):\n        popular_items = group['article_id'].value_counts().head(50).to_dict()\n        category_popular[category] = popular_items\n    \n    print(f\"    Generated category items for {len(category_popular)} categories\")\n    return category_popular\n\n\ndef calculate_seasonal_trending(transactions_df, cutoff_date, lookback_days=7):\n    \"\"\"Phase 1 Strategy 5: Seasonal trending items.\"\"\"\n    print(\"  [Strategy 5] Calculating seasonal trending items...\")\n    lookback_date = cutoff_date - timedelta(days=lookback_days)\n    recent_transactions = transactions_df[\n        (transactions_df['t_dat'] >= lookback_date) &\n        (transactions_df['t_dat'] < cutoff_date)\n    ].copy()\n    \n    # Calculate trending items (recent popularity surge)\n    trending_items = recent_transactions['article_id'].value_counts().head(100).to_dict()\n    \n    print(f\"    Generated {len(trending_items)} trending items\")\n    return trending_items\n\n\ndef get_product_code_prefix(article_id, articles_df):\n    \"\"\"Get product code prefix for diversity penalty.\"\"\"\n    if 'product_code' in articles_df.columns:\n        product_code = articles_df[articles_df['article_id'] == article_id]['product_code'].values\n        if len(product_code) > 0 and pd.notna(product_code[0]):\n            # Use first 3-4 characters as prefix\n            return str(product_code[0])[:4]\n    return None\n\n\ndef generate_customer_recommendations(\n    customer_id, \n    transactions_df, \n    articles_df, \n    customers_df,\n    cutoff_date,\n    repurchase_candidates,\n    age_group_popular,\n    cooccurrence_items,\n    category_items,\n    seasonal_trending,\n    n=12\n):\n    \"\"\"Generate recommendations for a single customer using all 5 strategies.\"\"\"\n    candidate_scores = defaultdict(float)\n    \n    # Get customer info\n    customer_info = customers_df[customers_df['customer_id'] == customer_id]\n    age_group = None\n    if len(customer_info) > 0 and 'age' in customer_info.columns:\n        age = customer_info['age'].values[0]\n        if pd.notna(age):\n            if age <= 25:\n                age_group = '18-25'\n            elif age <= 35:\n                age_group = '26-35'\n            elif age <= 45:\n                age_group = '36-45'\n            elif age <= 55:\n                age_group = '46-55'\n            else:\n                age_group = '55+'\n    \n    # Strategy 1: Repurchase candidates\n    if customer_id in repurchase_candidates:\n        for article_id, purchase_count in repurchase_candidates[customer_id].items():\n            # Get purchase history for this article\n            customer_history = transactions_df[\n                (transactions_df['customer_id'] == customer_id) &\n                (transactions_df['article_id'] == article_id) &\n                (transactions_df['t_dat'] < cutoff_date)\n            ]\n            \n            if len(customer_history) > 0:\n                # Calculate weighted score\n                total_score = 0.0\n                for _, row in customer_history.iterrows():\n                    days_ago = (cutoff_date - row['t_dat']).days\n                    time_weight = exponential_time_decay(days_ago)\n                    timescale_mult = get_timescale_multiplier(days_ago)\n                    recency_mult = get_recency_multiplier(days_ago)\n                    \n                    score = time_weight * timescale_mult * recency_mult\n                    total_score += score\n                \n                # Phase 4: Purchase frequency boost\n                freq_boost = 1 + 0.2 * purchase_count\n                candidate_scores[article_id] += total_score * freq_boost * 1.0  # Strategy weight\n    \n    # Strategy 2: Age group popular items\n    if age_group and age_group in age_group_popular:\n        for article_id, popularity in age_group_popular[age_group].items():\n            candidate_scores[article_id] += popularity * 0.8  # Strategy weight\n    \n    # Strategy 3: Co-occurrence items\n    customer_articles = transactions_df[\n        (transactions_df['customer_id'] == customer_id) &\n        (transactions_df['t_dat'] < cutoff_date)\n    ]['article_id'].unique()\n    \n    for article_id in customer_articles:\n        if article_id in cooccurrence_items:\n            for cooccur_article, count in cooccurrence_items[article_id].items():\n                candidate_scores[cooccur_article] += count * 1.2  # Strategy weight\n    \n    # Strategy 4: Same category items\n    customer_articles_list = transactions_df[\n        (transactions_df['customer_id'] == customer_id) &\n        (transactions_df['t_dat'] < cutoff_date)\n    ]['article_id'].tolist()\n    \n    if len(customer_articles_list) > 0:\n        # Get categories of customer's articles\n        customer_articles_df = pd.DataFrame({'article_id': customer_articles_list})\n        merged = customer_articles_df.merge(\n            articles_df[['article_id', 'product_type_no']] if 'product_type_no' in articles_df.columns \n            else articles_df[['article_id']],\n            on='article_id',\n            how='left'\n        )\n        \n        if 'product_type_no' in merged.columns:\n            customer_categories = merged['product_type_no'].dropna().unique()\n            for category in customer_categories:\n                if category in category_items:\n                    for article_id, popularity in category_items[category].items():\n                        candidate_scores[article_id] += popularity * 0.6  # Strategy weight\n    \n    # Strategy 5: Seasonal trending\n    for article_id, popularity in seasonal_trending.items():\n        candidate_scores[article_id] += popularity * 0.5  # Strategy weight\n    \n    # Phase 5: Apply diversity penalty\n    if len(candidate_scores) > 0:\n        # Sort by score\n        sorted_candidates = sorted(candidate_scores.items(), key=lambda x: x[1], reverse=True)\n        \n        # Apply diversity penalty\n        final_recommendations = []\n        used_prefixes = set()\n        \n        for article_id, score in sorted_candidates:\n            if len(final_recommendations) >= n:\n                break\n            \n            # Check diversity\n            prefix = get_product_code_prefix(article_id, articles_df)\n            if prefix and prefix in used_prefixes:\n                score *= (1 - DIVERSITY_PENALTY)  # Apply 30% penalty\n            \n            final_recommendations.append((article_id, score))\n            if prefix:\n                used_prefixes.add(prefix)\n        \n        # Sort again after penalty and return top n\n        final_recommendations.sort(key=lambda x: x[1], reverse=True)\n        recommendations = [article_id for article_id, _ in final_recommendations[:n]]\n        \n        # Fill with trending items if needed\n        if len(recommendations) < n:\n            for article_id in seasonal_trending.keys():\n                if article_id not in recommendations:\n                    recommendations.append(article_id)\n                    if len(recommendations) >= n:\n                        break\n        \n        return recommendations[:n]\n    \n    # Fallback: return trending items\n    return list(seasonal_trending.keys())[:n]\n\n\ndef process_customer_batch(args):\n    \"\"\"Worker function for multiprocessing - process a batch of customers.\"\"\"\n    (\n        customer_batch,\n        transactions_df,\n        articles_df,\n        customers_df,\n        cutoff_date,\n        repurchase_candidates,\n        age_group_popular,\n        cooccurrence_items,\n        category_items,\n        seasonal_trending,\n        n\n    ) = args\n\n    batch_recommendations = {}\n    for customer_id in customer_batch:\n        recs = generate_customer_recommendations(\n            customer_id,\n            transactions_df,\n            articles_df,\n            customers_df,\n            cutoff_date,\n            repurchase_candidates,\n            age_group_popular,\n            cooccurrence_items,\n            category_items,\n            seasonal_trending,\n            n=n\n        )\n        batch_recommendations[customer_id] = recs\n\n    return batch_recommendations\n\n\ndef generate_all_recommendations(\n    transactions_df,\n    articles_df,\n    customers_df,\n    customer_ids,\n    cutoff_date,\n    repurchase_candidates,\n    age_group_popular,\n    cooccurrence_items,\n    category_items,\n    seasonal_trending,\n    n=12,\n    batch_size=5000,\n    use_multiprocessing=True\n):\n    \"\"\"Generate recommendations for all customers using multiprocessing.\"\"\"\n    recommendations = {}\n    total_customers = len(customer_ids)\n\n    print(f\"Generating recommendations for {total_customers} customers...\")\n\n    if use_multiprocessing and total_customers > batch_size:\n        # Use multiprocessing for parallel processing\n        print(f\"  Using multiprocessing with batch_size={batch_size}\")\n\n        # Prepare arguments for multiprocessing\n        customer_batches = [customer_ids[i:i+batch_size] for i in range(0, total_customers, batch_size)]\n        num_batches = len(customer_batches)\n\n        # Prepare arguments for each batch\n        mp_args = [(\n            batch,\n            transactions_df,\n            articles_df,\n            customers_df,\n            cutoff_date,\n            repurchase_candidates,\n            age_group_popular,\n            cooccurrence_items,\n            category_items,\n            seasonal_trending,\n            n\n        ) for batch in customer_batches]\n\n        # Use multiprocessing pool\n        num_processes = min(mp.cpu_count(), num_batches)\n        print(f\"  Using {num_processes} processes...\")\n\n        with mp.Pool(processes=num_processes) as pool:\n            # Process batches in parallel with progress bar\n            results = []\n            for result in tqdm(\n                pool.imap(process_customer_batch, mp_args),\n                total=num_batches,\n                desc=\"Processing batches (parallel)\"\n            ):\n                results.append(result)\n\n        # Combine results\n        for batch_result in results:\n            recommendations.update(batch_result)\n\n    else:\n        # Fallback to sequential processing\n        print(\"  Using sequential processing...\")\n        for i in tqdm(range(0, total_customers, batch_size), desc=\"Processing batches\"):\n            batch = customer_ids[i:i+batch_size]\n\n            for customer_id in batch:\n                recs = generate_customer_recommendations(\n                    customer_id,\n                    transactions_df,\n                    articles_df,\n                    customers_df,\n                    cutoff_date,\n                    repurchase_candidates,\n                    age_group_popular,\n                    cooccurrence_items,\n                    category_items,\n                    seasonal_trending,\n                    n=n\n                )\n                recommendations[customer_id] = recs\n\n    print(f\"  Completed processing {total_customers} customers\")\n    return recommendations\n\n\ndef create_submission(recommendations, output_file='submission.csv'):\n    \"\"\"Create submission file in required format.\"\"\"\n    submission_data = []\n    \n    for customer_id, article_list in recommendations.items():\n        prediction = ' '.join([str(article_id).zfill(10) for article_id in article_list])\n        submission_data.append({\n            'customer_id': customer_id,\n            'prediction': prediction\n        })\n    \n    submission_df = pd.DataFrame(submission_data)\n    submission_df.to_csv(output_file, index=False)\n    print(f\"Submission saved to {output_file}\")\n    print(f\"Total predictions: {len(submission_df)}\")\n    \n    return submission_df\n\n\ndef main():\n    \"\"\"Main execution function.\"\"\"\n    start_time = datetime.now()\n    print(\"=\"*80)\n    print(\"H&M Personalized Fashion Recommendations - V2.8 Enhanced\")\n    print(\"=\"*80)\n    print_memory_usage(\"start\")\n\n    # 1. Load Data\n    print(\"\\n[1] Loading data...\")\n    try:\n        # Chunked loading for transactions if file is very large\n        transactions_path = os.path.join(DATA_DIR, \"transactions_train.csv\")\n        if os.path.exists(transactions_path):\n            file_size_gb = os.path.getsize(transactions_path) / (1024**3)\n            print(f\"  Transactions file size: {file_size_gb:.2f} GB\")\n\n            if file_size_gb > 2.0:  # Use chunked loading for large files\n                print(\"  Using chunked loading for large transactions file...\")\n                chunks = []\n                chunk_size = 1000000  # 1M rows per chunk\n                for chunk in pd.read_csv(\n                    transactions_path,\n                    parse_dates=['t_dat'],\n                    chunksize=chunk_size\n                ):\n                    chunks.append(chunk)\n                    print(f\"    Loaded chunk with {len(chunk):,} rows\")\n                    if PSUTIL_AVAILABLE:\n                        print_memory_usage(\"after chunk\")\n\n                transactions = pd.concat(chunks, ignore_index=True)\n                del chunks\n                gc.collect()\n            else:\n                transactions = pd.read_csv(transactions_path, parse_dates=['t_dat'])\n        else:\n            raise FileNotFoundError(f\"Transactions file not found: {transactions_path}\")\n\n        articles = pd.read_csv(os.path.join(DATA_DIR, \"articles.csv\"))\n        customers = pd.read_csv(os.path.join(DATA_DIR, \"customers.csv\"))\n        sample_submission = pd.read_csv(os.path.join(DATA_DIR, \"sample_submission.csv\"))\n\n        print(f\"  Transactions: {len(transactions):,} rows\")\n        print(f\"  Articles: {len(articles):,} rows\")\n        print(f\"  Customers: {len(customers):,} rows\")\n        print(f\"  Sample submission: {len(sample_submission):,} rows\")\n        print_memory_usage(\"after loading\")\n\n    except Exception as e:\n        print(f\"  Error loading data: {e}\")\n        sys.exit(1)\n    \n    # Convert IDs to strings\n    transactions['customer_id'] = transactions['customer_id'].astype(str)\n    transactions['article_id'] = transactions['article_id'].astype(str)\n    articles['article_id'] = articles['article_id'].astype(str)\n    customers['customer_id'] = customers['customer_id'].astype(str)\n    sample_submission['customer_id'] = sample_submission['customer_id'].astype(str)\n    \n    # 2. Memory Optimization\n    print(\"\\n[2] Optimizing memory usage...\")\n    transactions = reduce_mem_usage(transactions)\n    articles = reduce_mem_usage(articles)\n    customers = reduce_mem_usage(customers)\n    gc.collect()\n    \n    # 3. Determine cutoff date\n    print(\"\\n[3] Setting up prediction date...\")\n    cutoff_date = transactions['t_dat'].max()\n    print(f\"  Cutoff date: {cutoff_date}\")\n    print(f\"  Date range: {transactions['t_dat'].min()} to {cutoff_date}\")\n    \n    # 4. Phase 1: Generate candidates using 5 strategies\n    print(\"\\n[4] Phase 1: Enhanced Candidate Generation (5 strategies)...\")\n    print_memory_usage(\"before candidate generation\")\n\n    repurchase_candidates = calculate_repurchase_candidates(transactions, cutoff_date, lookback_days=90)\n    print_memory_usage(\"after repurchase candidates\")\n    gc.collect()\n\n    age_group_popular = calculate_age_group_popular_items(transactions, customers, cutoff_date, lookback_days=30)\n    print_memory_usage(\"after age group popular\")\n    gc.collect()\n\n    cooccurrence_items = calculate_cooccurrence_items(transactions, cutoff_date, lookback_days=90, min_cooccurrence=5)\n    print_memory_usage(\"after co-occurrence\")\n    gc.collect()\n\n    category_items = calculate_category_items(transactions, articles, cutoff_date, lookback_days=90)\n    print_memory_usage(\"after category items\")\n    gc.collect()\n\n    seasonal_trending = calculate_seasonal_trending(transactions, cutoff_date, lookback_days=7)\n    print_memory_usage(\"after seasonal trending\")\n    gc.collect()\n\n    # 5. Generate recommendations\n    print(\"\\n[5] Generating recommendations for all customers...\")\n    test_customers = sample_submission['customer_id'].unique()\n    \n    recommendations = generate_all_recommendations(\n        transactions,\n        articles,\n        customers,\n        test_customers,\n        cutoff_date,\n        repurchase_candidates,\n        age_group_popular,\n        cooccurrence_items,\n        category_items,\n        seasonal_trending,\n        n=12,\n        batch_size=5000\n    )\n\n    # Memory cleanup - remove large dataframes we don't need anymore\n    del transactions\n    gc.collect()\n    print_memory_usage(\"after cleanup\")\n\n    # 6. Create submission\n    print(\"\\n[6] Creating submission file...\")\n    submission_df = create_submission(recommendations, output_file='submission.csv')\n    \n    # 7. Summary\n    end_time = datetime.now()\n    runtime = (end_time - start_time).total_seconds() / 60\n    final_memory = get_memory_usage()\n\n    print(\"\\n\" + \"=\"*80)\n    print(\"V2.8 ENHANCEMENT COMPLETE\")\n    print(\"=\"*80)\n    print(f\"Total runtime: {runtime:.2f} minutes\")\n    print(f\"Target met: {'✅' if 27 <= runtime <= 39 else '❌'} (27-39 min target)\")\n    print(f\"Final memory usage: {final_memory:.1f} MB\")\n    print(f\"Total predictions: {len(submission_df):,}\")\n    print(f\"Submission file: submission.csv\")\n    print(\"\\nOPTIMIZATIONS APPLIED:\")\n    print(f\"  • Vectorized co-occurrence calculation\")\n    print(f\"  • Multiprocessing for customer batches ({mp.cpu_count()} cores)\")\n    print(f\"  • Dask parallel processing for large operations\")\n    print(f\"  • Memory optimization and chunked loading\")\n    print(f\"  • Aggressive garbage collection\")\n    print(\"=\"*80)\n\n    # Final cleanup\n    gc.collect()\n\n    return submission_df\n\n\nif __name__ == \"__main__\":\n    # Required for multiprocessing on Windows\n    mp.freeze_support()\n    submission_df = main()","metadata":{"trusted":true,"execution":{"iopub.status.busy":"2025-11-07T15:28:44.426755Z","iopub.execute_input":"2025-11-07T15:28:44.427276Z","execution_failed":"2025-11-07T15:44:02.198Z"}},"outputs":[],"execution_count":null},{"cell_type":"code","source":"\"\"\"\nH&M Personalized Fashion Recommendations - V2.8 Enhanced Implementation\nFollowing the V2 improvement plan for enhanced recommendation system\nExpected MAP@12: 0.032-0.038 (vs 0.012-0.015 baseline)\nRuntime: 27-39 minutes (target: <1 hour)\nOptimized for Kaggle GPU environment\n\"\"\"\n\nimport pandas as pd\nimport numpy as np\nfrom datetime import datetime, timedelta\nfrom collections import defaultdict, Counter\nimport warnings\nimport os\nimport sys\nimport gc\nfrom tqdm import tqdm\nimport multiprocessing as mp\nfrom itertools import combinations\n\n# Try to import optional dependencies\ntry:\n    import psutil\n    PSUTIL_AVAILABLE = True\nexcept ImportError:\n    PSUTIL_AVAILABLE = False\n\n# Try to import dask for parallel dataframe operations\ntry:\n    import dask.dataframe as dd\n    from dask import delayed, compute\n    DASK_AVAILABLE = True\nexcept ImportError:\n    print(\"Warning: Dask not available, using pandas only\")\n    DASK_AVAILABLE = False\n\n# Memory monitoring\ndef get_memory_usage():\n    \"\"\"Get current memory usage in MB.\"\"\"\n    if PSUTIL_AVAILABLE:\n        process = psutil.Process()\n        return process.memory_info().rss / 1024 / 1024\n    else:\n        return 0.0\n\ndef print_memory_usage(message=\"\"):\n    \"\"\"Print current memory usage.\"\"\"\n    if PSUTIL_AVAILABLE:\n        mem_mb = get_memory_usage()\n        mem_gb = mem_mb / 1024\n        warning = \" ⚠️ HIGH MEMORY\" if mem_gb > MAX_MEMORY_GB else \"\"\n        print(f\"    Memory usage{ ' - ' + message if message else ''}: {mem_mb:.1f} MB ({mem_gb:.1f} GB){warning}\")\n\ndef check_memory_limit():\n    \"\"\"Check if memory usage exceeds limit and force garbage collection.\"\"\"\n    if PSUTIL_AVAILABLE:\n        mem_gb = get_memory_usage() / 1024\n        if mem_gb > MAX_MEMORY_GB:\n            print(f\"    ⚠️ Memory usage ({mem_gb:.1f} GB) exceeds limit ({MAX_MEMORY_GB} GB), forcing GC...\")\n            gc.collect()\n            mem_gb_after = get_memory_usage() / 1024\n            print(f\"    Memory after GC: {mem_gb_after:.1f} GB\")\n            return mem_gb_after > MAX_MEMORY_GB\n    return False\n\nwarnings.filterwarnings('ignore')\n\n# Constants\nDATA_DIR = \"/kaggle/input/h-and-m-personalized-fashion-recommendations\"\nLAMBDA_DECAY = np.log(2) / 14  # Exponential decay parameter (14-day half-life)\nDIVERSITY_PENALTY = 0.3  # 30% reduction for same product code prefix\n\n# Memory optimization settings for Kaggle\nCHUNK_SIZE = 500000  # Smaller chunks to reduce memory usage\nMAX_MEMORY_GB = 13  # Target memory limit (increased for Kaggle GPU)\nDISABLE_MULTIPROCESSING = False  # Keep multiprocessing but optimize it\nDISABLE_DASK = True  # Disable Dask to save memory\nMAX_WORKERS = 2  # Limit workers to avoid memory spikes\nREC_BATCH_SIZE = 1000  # Smaller batch size for recommendations\n\n\ndef reduce_mem_usage(df):\n    \"\"\"Reduce memory usage by downcasting numeric dtypes.\"\"\"\n    start_mem = df.memory_usage().sum() / 1024**2\n\n    for col in df.columns:\n        col_type = df[col].dtype\n\n        # Skip datetime and categorical columns\n        if col_type == 'datetime64[ns]' or col_type.name == 'category':\n            continue\n\n        if col_type != object:\n            c_min = df[col].min()\n            c_max = df[col].max()\n\n            if str(col_type)[:3] == 'int':\n                if c_min > np.iinfo(np.int8).min and c_max < np.iinfo(np.int8).max:\n                    df[col] = df[col].astype(np.int8)\n                elif c_min > np.iinfo(np.int16).min and c_max < np.iinfo(np.int16).max:\n                    df[col] = df[col].astype(np.int16)\n                elif c_min > np.iinfo(np.int32).min and c_max < np.iinfo(np.int32).max:\n                    df[col] = df[col].astype(np.int32)\n            elif str(col_type)[:5] == 'float':\n                # Skip float16 downcasting due to pandas compatibility issues\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.float64)\n\n    end_mem = df.memory_usage().sum() / 1024**2\n    print(f'    Memory usage after optimization: {end_mem:.2f} MB')\n    print(f'    Decreased by {100 * (start_mem - end_mem) / start_mem:.1f}%')\n    return df\n\ndef exponential_time_decay(days_ago):\n    \"\"\"Phase 2: Exponential time decay with 14-day half-life.\"\"\"\n    return np.exp(-LAMBDA_DECAY * days_ago)\n\n\ndef get_timescale_multiplier(days_ago):\n    \"\"\"Phase 3: Multi-timescale windows.\"\"\"\n    if days_ago <= 7:\n        return 1.5  # Short-term: trending preferences\n    elif days_ago <= 30:\n        return 1.0  # Medium-term: stable preferences\n    elif days_ago <= 90:\n        return 0.5  # Long-term: general style\n    else:\n        return 0.1  # Very old: minimal weight\n\n\ndef get_recency_multiplier(days_ago):\n    \"\"\"Phase 4: Recency multiplier for purchase frequency boost.\"\"\"\n    if days_ago <= 3:\n        return 2.0  # Last 3 days: 2.0x\n    elif days_ago <= 7:\n        return 1.5  # Last 7 days: 1.5x\n    else:\n        return 1.0\n\n\ndef calculate_repurchase_candidates(transactions_df, cutoff_date, lookback_days=90):\n    \"\"\"Phase 1 Strategy 1: Repurchase candidates from customer history.\"\"\"\n    print(\"  [Strategy 1] Calculating repurchase candidates...\")\n    lookback_date = cutoff_date - timedelta(days=lookback_days)\n    recent_transactions = transactions_df[\n        (transactions_df['t_dat'] >= lookback_date) &\n        (transactions_df['t_dat'] < cutoff_date)\n    ].copy()\n\n    print(f\"    Processing {len(recent_transactions):,} transactions...\")\n\n    if DASK_AVAILABLE and len(recent_transactions) > 1000000 and not DISABLE_DASK:  # Use dask for large datasets\n        print(\"    Using Dask for parallel processing...\")\n\n        # Convert to dask dataframe\n        ddf = dd.from_pandas(recent_transactions, npartitions=mp.cpu_count())\n\n        # Parallel groupby operation\n        result = ddf.groupby(['customer_id', 'article_id']).size().compute()\n\n        # Convert to dictionary format\n        repurchase_candidates = {}\n        for (customer_id, article_id), count in result.items():\n            if customer_id not in repurchase_candidates:\n                repurchase_candidates[customer_id] = {}\n            repurchase_candidates[customer_id][article_id] = count\n\n    else:\n        # Use pandas for smaller datasets or when dask not available\n        print(\"    Using pandas processing...\")\n        repurchase_candidates = {}\n        for customer_id, group in tqdm(recent_transactions.groupby('customer_id'), desc=\"    Processing customers\"):\n            article_counts = group['article_id'].value_counts()\n            repurchase_candidates[customer_id] = article_counts.to_dict()\n\n    print(f\"    Generated repurchase candidates for {len(repurchase_candidates)} customers\")\n    return repurchase_candidates\n\n\ndef calculate_age_group_popular_items(transactions_df, customers_df, cutoff_date, lookback_days=30):\n    \"\"\"Phase 1 Strategy 2: Age group popular items.\"\"\"\n    print(\"  [Strategy 2] Calculating age group popular items...\")\n    lookback_date = cutoff_date - timedelta(days=lookback_days)\n    recent_transactions = transactions_df[\n        (transactions_df['t_dat'] >= lookback_date) &\n        (transactions_df['t_dat'] < cutoff_date)\n    ].copy()\n\n    # Merge with customer age\n    if 'age' not in customers_df.columns:\n        print(\"    Warning: Age column not found, skipping age group strategy\")\n        return {}\n\n    # Create age groups\n    customers_df = customers_df.copy()\n    customers_df['age_group'] = pd.cut(\n        customers_df['age'],\n        bins=[0, 25, 35, 45, 55, 100],\n        labels=['18-25', '26-35', '36-45', '46-55', '55+']\n    )\n\n    print(f\"    Processing {len(recent_transactions):,} transactions for age groups...\")\n\n    if DASK_AVAILABLE and len(recent_transactions) > 500000 and not DISABLE_DASK:  # Use dask for large merges\n        print(\"    Using Dask for parallel merge/groupby...\")\n\n        # Convert to dask dataframes\n        trans_ddf = dd.from_pandas(recent_transactions, npartitions=mp.cpu_count())\n        cust_ddf = dd.from_pandas(customers_df[['customer_id', 'age_group']], npartitions=4)\n\n        # Parallel merge\n        merged_ddf = trans_ddf.merge(cust_ddf, on='customer_id', how='left')\n\n        # Parallel groupby and value_counts equivalent\n        result = merged_ddf.groupby(['age_group', 'article_id']).size().compute()\n\n        # Convert to dictionary format\n        age_group_popular = {}\n        for (age_group, article_id), count in result.items():\n            if pd.notna(age_group):\n                if age_group not in age_group_popular:\n                    age_group_popular[age_group] = {}\n                age_group_popular[age_group][article_id] = count\n\n        # Keep only top 50 per age group\n        for age_group in age_group_popular:\n            sorted_items = sorted(age_group_popular[age_group].items(), key=lambda x: x[1], reverse=True)\n            age_group_popular[age_group] = dict(sorted_items[:50])\n\n    else:\n        # Use pandas for smaller datasets or when dask not available\n        print(\"    Using pandas processing...\")\n        merged = recent_transactions.merge(\n            customers_df[['customer_id', 'age_group']],\n            on='customer_id',\n            how='left'\n        )\n\n        # Calculate popular items per age group\n        age_group_popular = {}\n        for age_group in merged['age_group'].dropna().unique():\n            age_transactions = merged[merged['age_group'] == age_group]\n            popular_items = age_transactions['article_id'].value_counts().head(50).to_dict()\n            age_group_popular[age_group] = popular_items\n\n    print(f\"    Generated popular items for {len(age_group_popular)} age groups\")\n    return age_group_popular\n\n\ndef calculate_cooccurrence_items(transactions_df, cutoff_date, lookback_days=90, min_cooccurrence=5):\n    \"\"\"Phase 1 Strategy 3: Items bought together (co-occurrence) - VECTORIZED VERSION.\"\"\"\n    print(\"  [Strategy 3] Calculating co-occurrence patterns (vectorized)...\")\n    lookback_date = cutoff_date - timedelta(days=lookback_days)\n    recent_transactions = transactions_df[\n        (transactions_df['t_dat'] >= lookback_date) &\n        (transactions_df['t_dat'] < cutoff_date)\n    ].copy()\n\n    print(f\"    Processing {len(recent_transactions):,} transactions for co-occurrence...\")\n\n    # Vectorized approach: Group by customer and date, then create combinations\n    grouped = recent_transactions.groupby(['customer_id', 't_dat'])['article_id'].unique()\n    grouped = grouped.reset_index()\n\n    # Create all combinations for each customer-date group\n    all_combinations = []\n    for articles in tqdm(grouped['article_id'], desc=\"    Creating combinations\"):\n        if len(articles) > 1:\n            # Create unordered pairs (combinations, not permutations)\n            for art1, art2 in combinations(sorted(articles), 2):\n                all_combinations.append((art1, art2))\n\n    if not all_combinations:\n        print(\"    No co-occurrence patterns found\")\n        return {}\n\n    # Convert to DataFrame for efficient counting\n    combo_df = pd.DataFrame(all_combinations, columns=['article1', 'article2'])\n    cooccurrence_counts = combo_df.groupby(['article1', 'article2']).size().reset_index(name='count')\n\n    # Filter by minimum co-occurrence\n    cooccurrence_counts = cooccurrence_counts[cooccurrence_counts['count'] >= min_cooccurrence]\n\n    # Convert to dictionary format\n    cooccurrence = defaultdict(Counter)\n    for _, row in cooccurrence_counts.iterrows():\n        art1, art2, count = row['article1'], row['article2'], row['count']\n        cooccurrence[art1][art2] = count\n        cooccurrence[art2][art1] = count\n\n    # Filter out articles with no co-occurrences\n    filtered_cooccurrence = {k: dict(v) for k, v in cooccurrence.items() if v}\n\n    print(f\"    Generated co-occurrence patterns for {len(filtered_cooccurrence)} articles\")\n    return filtered_cooccurrence\n\n\ndef calculate_category_items(transactions_df, articles_df, cutoff_date, lookback_days=90):\n    \"\"\"Phase 1 Strategy 4: Same category items.\"\"\"\n    print(\"  [Strategy 4] Calculating category-based items...\")\n    lookback_date = cutoff_date - timedelta(days=lookback_days)\n    recent_transactions = transactions_df[\n        (transactions_df['t_dat'] >= lookback_date) &\n        (transactions_df['t_dat'] < cutoff_date)\n    ].copy()\n    \n    # Merge with article categories\n    if 'product_type_no' not in articles_df.columns:\n        print(\"    Warning: product_type_no not found, using article_id grouping\")\n        category_col = 'article_id'\n    else:\n        category_col = 'product_type_no'\n    \n    merged = recent_transactions.merge(\n        articles_df[['article_id', category_col]], \n        on='article_id', \n        how='left'\n    )\n    \n    # Calculate popular items per category\n    category_popular = {}\n    for category, group in merged.groupby(category_col):\n        popular_items = group['article_id'].value_counts().head(50).to_dict()\n        category_popular[category] = popular_items\n    \n    print(f\"    Generated category items for {len(category_popular)} categories\")\n    return category_popular\n\n\ndef calculate_seasonal_trending(transactions_df, cutoff_date, lookback_days=7):\n    \"\"\"Phase 1 Strategy 5: Seasonal trending items.\"\"\"\n    print(\"  [Strategy 5] Calculating seasonal trending items...\")\n    lookback_date = cutoff_date - timedelta(days=lookback_days)\n    recent_transactions = transactions_df[\n        (transactions_df['t_dat'] >= lookback_date) &\n        (transactions_df['t_dat'] < cutoff_date)\n    ].copy()\n    \n    # Calculate trending items (recent popularity surge)\n    trending_items = recent_transactions['article_id'].value_counts().head(100).to_dict()\n    \n    print(f\"    Generated {len(trending_items)} trending items\")\n    return trending_items\n\n\ndef get_product_code_prefix(article_id, articles_df):\n    \"\"\"Get product code prefix for diversity penalty.\"\"\"\n    if 'product_code' in articles_df.columns:\n        product_code = articles_df[articles_df['article_id'] == article_id]['product_code'].values\n        if len(product_code) > 0 and pd.notna(product_code[0]):\n            # Use first 3-4 characters as prefix\n            return str(product_code[0])[:4]\n    return None\n\n\ndef generate_customer_recommendations(\n    customer_id, \n    transactions_df, \n    articles_df, \n    customers_df,\n    cutoff_date,\n    repurchase_candidates,\n    age_group_popular,\n    cooccurrence_items,\n    category_items,\n    seasonal_trending,\n    n=12\n):\n    \"\"\"Generate recommendations for a single customer using all 5 strategies.\"\"\"\n    candidate_scores = defaultdict(float)\n    \n    # Get customer info\n    customer_info = customers_df[customers_df['customer_id'] == customer_id]\n    age_group = None\n    if len(customer_info) > 0 and 'age' in customer_info.columns:\n        age = customer_info['age'].values[0]\n        if pd.notna(age):\n            if age <= 25:\n                age_group = '18-25'\n            elif age <= 35:\n                age_group = '26-35'\n            elif age <= 45:\n                age_group = '36-45'\n            elif age <= 55:\n                age_group = '46-55'\n            else:\n                age_group = '55+'\n    \n    # Strategy 1: Repurchase candidates\n    if customer_id in repurchase_candidates:\n        for article_id, purchase_count in repurchase_candidates[customer_id].items():\n            # Get purchase history for this article\n            customer_history = transactions_df[\n                (transactions_df['customer_id'] == customer_id) &\n                (transactions_df['article_id'] == article_id) &\n                (transactions_df['t_dat'] < cutoff_date)\n            ]\n            \n            if len(customer_history) > 0:\n                # Calculate weighted score\n                total_score = 0.0\n                for _, row in customer_history.iterrows():\n                    days_ago = (cutoff_date - row['t_dat']).days\n                    time_weight = exponential_time_decay(days_ago)\n                    timescale_mult = get_timescale_multiplier(days_ago)\n                    recency_mult = get_recency_multiplier(days_ago)\n                    \n                    score = time_weight * timescale_mult * recency_mult\n                    total_score += score\n                \n                # Phase 4: Purchase frequency boost\n                freq_boost = 1 + 0.2 * purchase_count\n                candidate_scores[article_id] += total_score * freq_boost * 1.0  # Strategy weight\n    \n    # Strategy 2: Age group popular items\n    if age_group and age_group in age_group_popular:\n        for article_id, popularity in age_group_popular[age_group].items():\n            candidate_scores[article_id] += popularity * 0.8  # Strategy weight\n    \n    # Strategy 3: Co-occurrence items\n    customer_articles = transactions_df[\n        (transactions_df['customer_id'] == customer_id) &\n        (transactions_df['t_dat'] < cutoff_date)\n    ]['article_id'].unique()\n    \n    for article_id in customer_articles:\n        if article_id in cooccurrence_items:\n            for cooccur_article, count in cooccurrence_items[article_id].items():\n                candidate_scores[cooccur_article] += count * 1.2  # Strategy weight\n    \n    # Strategy 4: Same category items\n    customer_articles_list = transactions_df[\n        (transactions_df['customer_id'] == customer_id) &\n        (transactions_df['t_dat'] < cutoff_date)\n    ]['article_id'].tolist()\n    \n    if len(customer_articles_list) > 0:\n        # Get categories of customer's articles\n        customer_articles_df = pd.DataFrame({'article_id': customer_articles_list})\n        merged = customer_articles_df.merge(\n            articles_df[['article_id', 'product_type_no']] if 'product_type_no' in articles_df.columns \n            else articles_df[['article_id']],\n            on='article_id',\n            how='left'\n        )\n        \n        if 'product_type_no' in merged.columns:\n            customer_categories = merged['product_type_no'].dropna().unique()\n            for category in customer_categories:\n                if category in category_items:\n                    for article_id, popularity in category_items[category].items():\n                        candidate_scores[article_id] += popularity * 0.6  # Strategy weight\n    \n    # Strategy 5: Seasonal trending\n    for article_id, popularity in seasonal_trending.items():\n        candidate_scores[article_id] += popularity * 0.5  # Strategy weight\n    \n    # Phase 5: Apply diversity penalty\n    if len(candidate_scores) > 0:\n        # Sort by score\n        sorted_candidates = sorted(candidate_scores.items(), key=lambda x: x[1], reverse=True)\n        \n        # Apply diversity penalty\n        final_recommendations = []\n        used_prefixes = set()\n        \n        for article_id, score in sorted_candidates:\n            if len(final_recommendations) >= n:\n                break\n            \n            # Check diversity\n            prefix = get_product_code_prefix(article_id, articles_df)\n            if prefix and prefix in used_prefixes:\n                score *= (1 - DIVERSITY_PENALTY)  # Apply 30% penalty\n            \n            final_recommendations.append((article_id, score))\n            if prefix:\n                used_prefixes.add(prefix)\n        \n        # Sort again after penalty and return top n\n        final_recommendations.sort(key=lambda x: x[1], reverse=True)\n        recommendations = [article_id for article_id, _ in final_recommendations[:n]]\n        \n        # Fill with trending items if needed\n        if len(recommendations) < n:\n            for article_id in seasonal_trending.keys():\n                if article_id not in recommendations:\n                    recommendations.append(article_id)\n                    if len(recommendations) >= n:\n                        break\n        \n        return recommendations[:n]\n    \n    # Fallback: return trending items\n    return list(seasonal_trending.keys())[:n]\n\n\ndef process_customer_batch(args):\n    \"\"\"Worker function for multiprocessing - process a batch of customers.\"\"\"\n    (\n        customer_batch,\n        transactions_df,\n        articles_df,\n        customers_df,\n        cutoff_date,\n        repurchase_candidates,\n        age_group_popular,\n        cooccurrence_items,\n        category_items,\n        seasonal_trending,\n        n\n    ) = args\n\n    batch_recommendations = {}\n    for customer_id in customer_batch:\n        recs = generate_customer_recommendations(\n            customer_id,\n            transactions_df,\n            articles_df,\n            customers_df,\n            cutoff_date,\n            repurchase_candidates,\n            age_group_popular,\n            cooccurrence_items,\n            category_items,\n            seasonal_trending,\n            n=n\n        )\n        batch_recommendations[customer_id] = recs\n\n    return batch_recommendations\n\n\ndef generate_all_recommendations(\n    transactions_df,\n    articles_df,\n    customers_df,\n    customer_ids,\n    cutoff_date,\n    repurchase_candidates,\n    age_group_popular,\n    cooccurrence_items,\n    category_items,\n    seasonal_trending,\n    n=12,\n    batch_size=None,\n    use_multiprocessing=True\n):\n    \"\"\"Generate recommendations for all customers using optimized multiprocessing.\"\"\"\n    if batch_size is None:\n        batch_size = REC_BATCH_SIZE\n    \n    recommendations = {}\n    total_customers = len(customer_ids)\n\n    print(f\"Generating recommendations for {total_customers} customers...\")\n\n    # Check memory before starting\n    if PSUTIL_AVAILABLE:\n        mem_gb = get_memory_usage() / 1024\n        print(f\"  Current memory: {mem_gb:.1f} GB\")\n        if mem_gb > MAX_MEMORY_GB * 0.9:  # If already at 90% of limit, use sequential\n            print(f\"  ⚠️ Memory usage high ({mem_gb:.1f} GB), using sequential processing\")\n            use_multiprocessing = False\n\n    if use_multiprocessing and total_customers > batch_size and not DISABLE_MULTIPROCESSING:\n        # Use multiprocessing for parallel processing\n        num_batches = (total_customers + batch_size - 1) // batch_size\n        print(f\"  Using multiprocessing with batch_size={batch_size}\")\n        print(f\"  Processing {total_customers:,} customers in {num_batches} batches...\")\n\n        # Limit number of workers to avoid memory spikes\n        num_workers = min(MAX_WORKERS, num_batches, mp.cpu_count())\n        print(f\"  Using {num_workers} worker processes (limited to avoid memory issues)...\")\n\n        # Prepare arguments for multiprocessing\n        customer_batches = [customer_ids[i:i+batch_size] for i in range(0, total_customers, batch_size)]\n\n        # Prepare arguments for each batch\n        mp_args = [(\n            batch,\n            transactions_df,\n            articles_df,\n            customers_df,\n            cutoff_date,\n            repurchase_candidates,\n            age_group_popular,\n            cooccurrence_items,\n            category_items,\n            seasonal_trending,\n            n\n        ) for batch in customer_batches]\n\n        # Use multiprocessing pool with limited workers\n        with mp.Pool(processes=num_workers) as pool:\n            # Process batches in parallel with progress bar\n            results = []\n            batch_num = 0\n            for result in tqdm(\n                pool.imap(process_customer_batch, mp_args),\n                total=num_batches,\n                desc=\"Processing batches (parallel)\"\n            ):\n                results.append(result)\n                batch_num += 1\n                \n                # Periodic memory check and cleanup every 20 batches\n                if batch_num % 20 == 0:\n                    gc.collect()\n                    if PSUTIL_AVAILABLE:\n                        mem_gb = get_memory_usage() / 1024\n                        if mem_gb > MAX_MEMORY_GB:\n                            print(f\"    ⚠️ Memory at {mem_gb:.1f} GB after batch {batch_num}, forcing GC...\")\n                            gc.collect()\n\n        # Combine results\n        for batch_result in results:\n            recommendations.update(batch_result)\n\n    else:\n        # Fallback to sequential processing\n        num_batches = (total_customers + batch_size - 1) // batch_size\n        print(f\"  Using sequential processing with batch_size={batch_size}...\")\n        print(f\"  Processing {total_customers:,} customers in {num_batches} batches...\")\n        \n        for i in tqdm(range(0, total_customers, batch_size), desc=\"Processing batches\"):\n            batch = customer_ids[i:i+batch_size]\n\n            for customer_id in batch:\n                recs = generate_customer_recommendations(\n                    customer_id,\n                    transactions_df,\n                    articles_df,\n                    customers_df,\n                    cutoff_date,\n                    repurchase_candidates,\n                    age_group_popular,\n                    cooccurrence_items,\n                    category_items,\n                    seasonal_trending,\n                    n=n\n                )\n                recommendations[customer_id] = recs\n            \n            # Periodic memory check and cleanup\n            if (i // batch_size) % 10 == 0:\n                gc.collect()\n                if PSUTIL_AVAILABLE:\n                    mem_gb = get_memory_usage() / 1024\n                    if mem_gb > MAX_MEMORY_GB:\n                        print(f\"    ⚠️ Memory at {mem_gb:.1f} GB after batch {i//batch_size}, forcing GC...\")\n                        gc.collect()\n\n    print(f\"  Completed processing {total_customers} customers\")\n    return recommendations\n\n\ndef create_submission(recommendations, output_file='submission.csv'):\n    \"\"\"Create submission file in required format.\"\"\"\n    submission_data = []\n    \n    for customer_id, article_list in recommendations.items():\n        prediction = ' '.join([str(article_id).zfill(10) for article_id in article_list])\n        submission_data.append({\n            'customer_id': customer_id,\n            'prediction': prediction\n        })\n    \n    submission_df = pd.DataFrame(submission_data)\n    submission_df.to_csv(output_file, index=False)\n    print(f\"Submission saved to {output_file}\")\n    print(f\"Total predictions: {len(submission_df)}\")\n    \n    return submission_df\n\n\ndef main():\n    \"\"\"Main execution function.\"\"\"\n    start_time = datetime.now()\n    print(\"=\"*80)\n    print(\"H&M Personalized Fashion Recommendations - V2.8 Enhanced\")\n    print(\"=\"*80)\n    print_memory_usage(\"start\")\n\n    # 1. Load Data\n    print(\"\\n[1] Loading data...\")\n    try:\n        # Chunked loading for transactions if file is very large\n        transactions_path = os.path.join(DATA_DIR, \"transactions_train.csv\")\n        if os.path.exists(transactions_path):\n            file_size_gb = os.path.getsize(transactions_path) / (1024**3)\n            print(f\"  Transactions file size: {file_size_gb:.2f} GB\")\n\n            if file_size_gb > 1.0:  # Use chunked loading for large files\n                print(f\"  Using chunked loading for large transactions file (chunk_size={CHUNK_SIZE:,})...\")\n                chunks = []\n                for chunk in pd.read_csv(\n                    transactions_path,\n                    parse_dates=['t_dat'],\n                    chunksize=CHUNK_SIZE\n                ):\n                    chunks.append(chunk)\n                    print(f\"    Loaded chunk with {len(chunk):,} rows\")\n                    if PSUTIL_AVAILABLE:\n                        print_memory_usage(\"after chunk\")\n\n                transactions = pd.concat(chunks, ignore_index=True)\n                del chunks\n                gc.collect()\n            else:\n                transactions = pd.read_csv(transactions_path, parse_dates=['t_dat'])\n        else:\n            raise FileNotFoundError(f\"Transactions file not found: {transactions_path}\")\n\n        articles = pd.read_csv(os.path.join(DATA_DIR, \"articles.csv\"))\n        customers = pd.read_csv(os.path.join(DATA_DIR, \"customers.csv\"))\n        sample_submission = pd.read_csv(os.path.join(DATA_DIR, \"sample_submission.csv\"))\n\n        print(f\"  Transactions: {len(transactions):,} rows\")\n        print(f\"  Articles: {len(articles):,} rows\")\n        print(f\"  Customers: {len(customers):,} rows\")\n        print(f\"  Sample submission: {len(sample_submission):,} rows\")\n        print_memory_usage(\"after loading\")\n\n    except Exception as e:\n        print(f\"  Error loading data: {e}\")\n        sys.exit(1)\n    \n    # Convert IDs to strings\n    transactions['customer_id'] = transactions['customer_id'].astype(str)\n    transactions['article_id'] = transactions['article_id'].astype(str)\n    articles['article_id'] = articles['article_id'].astype(str)\n    customers['customer_id'] = customers['customer_id'].astype(str)\n    sample_submission['customer_id'] = sample_submission['customer_id'].astype(str)\n    \n    # 2. Memory Optimization\n    print(\"\\n[2] Optimizing memory usage...\")\n    transactions = reduce_mem_usage(transactions)\n    articles = reduce_mem_usage(articles)\n    customers = reduce_mem_usage(customers)\n    gc.collect()\n    \n    # 3. Determine cutoff date\n    print(\"\\n[3] Setting up prediction date...\")\n    cutoff_date = transactions['t_dat'].max()\n    print(f\"  Cutoff date: {cutoff_date}\")\n    print(f\"  Date range: {transactions['t_dat'].min()} to {cutoff_date}\")\n    \n    # 4. Phase 1: Generate candidates using 5 strategies\n    print(\"\\n[4] Phase 1: Enhanced Candidate Generation (5 strategies)...\")\n    print_memory_usage(\"before candidate generation\")\n\n    repurchase_candidates = calculate_repurchase_candidates(transactions, cutoff_date, lookback_days=90)\n    print_memory_usage(\"after repurchase candidates\")\n    if check_memory_limit():\n        print(\"    ⚠️ Continuing despite high memory usage...\")\n    gc.collect()\n\n    age_group_popular = calculate_age_group_popular_items(transactions, customers, cutoff_date, lookback_days=30)\n    print_memory_usage(\"after age group popular\")\n    if check_memory_limit():\n        print(\"    ⚠️ Continuing despite high memory usage...\")\n    gc.collect()\n\n    cooccurrence_items = calculate_cooccurrence_items(transactions, cutoff_date, lookback_days=90, min_cooccurrence=5)\n    print_memory_usage(\"after co-occurrence\")\n    if check_memory_limit():\n        print(\"    ⚠️ Continuing despite high memory usage...\")\n    gc.collect()\n\n    category_items = calculate_category_items(transactions, articles, cutoff_date, lookback_days=90)\n    print_memory_usage(\"after category items\")\n    if check_memory_limit():\n        print(\"    ⚠️ Continuing despite high memory usage...\")\n    gc.collect()\n\n    seasonal_trending = calculate_seasonal_trending(transactions, cutoff_date, lookback_days=7)\n    print_memory_usage(\"after seasonal trending\")\n    if check_memory_limit():\n        print(\"    ⚠️ Continuing despite high memory usage...\")\n    gc.collect()\n\n    # 5. Generate recommendations\n    print(\"\\n[5] Generating recommendations for all customers...\")\n    test_customers = sample_submission['customer_id'].unique()\n    \n    recommendations = generate_all_recommendations(\n        transactions,\n        articles,\n        customers,\n        test_customers,\n        cutoff_date,\n        repurchase_candidates,\n        age_group_popular,\n        cooccurrence_items,\n        category_items,\n        seasonal_trending,\n        n=12,\n        batch_size=REC_BATCH_SIZE,  # Use optimized batch size\n        use_multiprocessing=not DISABLE_MULTIPROCESSING\n    )\n\n    # Memory cleanup - remove large dataframes we don't need anymore\n    del transactions\n    gc.collect()\n    print_memory_usage(\"after cleanup\")\n\n    # 6. Create submission\n    print(\"\\n[6] Creating submission file...\")\n    submission_df = create_submission(recommendations, output_file='submission.csv')\n    \n    # 7. Summary\n    end_time = datetime.now()\n    runtime = (end_time - start_time).total_seconds() / 60\n    final_memory = get_memory_usage()\n\n    print(\"\\n\" + \"=\"*80)\n    print(\"V2.8 ENHANCEMENT COMPLETE\")\n    print(\"=\"*80)\n    print(f\"Total runtime: {runtime:.2f} minutes\")\n    print(f\"Target met: {'✅' if 27 <= runtime <= 39 else '❌'} (27-39 min target)\")\n    print(f\"Final memory usage: {final_memory:.1f} MB\")\n    print(f\"Total predictions: {len(submission_df):,}\")\n    print(f\"Submission file: submission.csv\")\n    print(\"\\nOPTIMIZATIONS APPLIED:\")\n    print(f\"  • Vectorized co-occurrence calculation\")\n    print(f\"  • Memory-optimized chunked loading ({CHUNK_SIZE:,} rows)\")\n    print(f\"  • Memory monitoring and limits ({MAX_MEMORY_GB}GB)\")\n    if not DISABLE_MULTIPROCESSING:\n        print(f\"  • Optimized multiprocessing ({MAX_WORKERS} workers, batch_size={REC_BATCH_SIZE})\")\n    else:\n        print(f\"  • Sequential processing (multiprocessing disabled for memory)\")\n    if not DISABLE_DASK:\n        print(f\"  • Dask parallel processing for large operations\")\n    else:\n        print(f\"  • Pandas processing (Dask disabled for memory)\")\n    print(f\"  • Aggressive garbage collection and cleanup\")\n    print(\"=\"*80)\n\n    # Final cleanup\n    gc.collect()\n\n    return submission_df\n\n\nif __name__ == \"__main__\":\n    # Required for multiprocessing on Windows\n    mp.freeze_support()\n    submission_df = main()","metadata":{"trusted":true,"execution":{"iopub.status.busy":"2025-11-07T15:47:51.986736Z","iopub.execute_input":"2025-11-07T15:47:51.987376Z"}},"outputs":[],"execution_count":null},{"cell_type":"code","source":"\"\"\"\nH&M Personalized Fashion Recommendations - V2.8 Enhanced Implementation\nFollowing the V2 improvement plan for enhanced recommendation system\nExpected MAP@12: 0.032-0.038 (vs 0.012-0.015 baseline)\nRuntime: 27-39 minutes (target: <1 hour)\nOptimized for Kaggle GPU environment\n\"\"\"\n\nimport pandas as pd\nimport numpy as np\nfrom datetime import datetime, timedelta\nfrom collections import defaultdict, Counter\nimport warnings\nimport os\nimport sys\nimport gc\nfrom tqdm import tqdm\nimport multiprocessing as mp\nfrom itertools import combinations\n\n# Try to import optional dependencies\ntry:\n    import psutil\n    PSUTIL_AVAILABLE = True\nexcept ImportError:\n    PSUTIL_AVAILABLE = False\n\n# Try to import dask for parallel dataframe operations\ntry:\n    import dask.dataframe as dd\n    from dask import delayed, compute\n    DASK_AVAILABLE = True\nexcept ImportError:\n    print(\"Warning: Dask not available, using pandas only\")\n    DASK_AVAILABLE = False\n\n# Memory monitoring\ndef get_memory_usage():\n    \"\"\"Get current memory usage in MB.\"\"\"\n    if PSUTIL_AVAILABLE:\n        process = psutil.Process()\n        return process.memory_info().rss / 1024 / 1024\n    else:\n        return 0.0\n\ndef print_memory_usage(message=\"\"):\n    \"\"\"Print current memory usage.\"\"\"\n    if PSUTIL_AVAILABLE:\n        mem_mb = get_memory_usage()\n        mem_gb = mem_mb / 1024\n        warning = \" ⚠️ HIGH MEMORY\" if mem_gb > MAX_MEMORY_GB else \"\"\n        print(f\"    Memory usage{ ' - ' + message if message else ''}: {mem_mb:.1f} MB ({mem_gb:.1f} GB){warning}\")\n\ndef check_memory_limit():\n    \"\"\"Check if memory usage exceeds limit and force garbage collection.\"\"\"\n    if PSUTIL_AVAILABLE:\n        mem_gb = get_memory_usage() / 1024\n        if mem_gb > MAX_MEMORY_GB:\n            print(f\"    ⚠️ Memory usage ({mem_gb:.1f} GB) exceeds limit ({MAX_MEMORY_GB} GB), forcing GC...\")\n            gc.collect()\n            mem_gb_after = get_memory_usage() / 1024\n            print(f\"    Memory after GC: {mem_gb_after:.1f} GB\")\n            return mem_gb_after > MAX_MEMORY_GB\n    return False\n\nwarnings.filterwarnings('ignore')\n\n# Constants\nDATA_DIR = \"/kaggle/input/h-and-m-personalized-fashion-recommendations\"\nLAMBDA_DECAY = np.log(2) / 14  # Exponential decay parameter (14-day half-life)\nDIVERSITY_PENALTY = 0.3  # 30% reduction for same product code prefix\n\n# Memory optimization settings for Kaggle\nCHUNK_SIZE = 500000  # Smaller chunks to reduce memory usage\nMAX_MEMORY_GB = 13  # Target memory limit (increased for Kaggle GPU)\nDISABLE_MULTIPROCESSING = True  # Disable multiprocessing to fix hang issue (Strategy 3)\nDISABLE_DASK = True  # Disable Dask to save memory\nMAX_WORKERS = 2  # Limit workers to avoid memory spikes\nREC_BATCH_SIZE = 1000  # Batch size for recommendations\nTRANSACTION_LOOKBACK_DAYS = 90  # Maximum lookback window for filtering\n\n\ndef reduce_mem_usage(df):\n    \"\"\"Reduce memory usage by downcasting numeric dtypes.\"\"\"\n    start_mem = df.memory_usage().sum() / 1024**2\n\n    for col in df.columns:\n        col_type = df[col].dtype\n\n        # Skip datetime and categorical columns\n        if col_type == 'datetime64[ns]' or col_type.name == 'category':\n            continue\n\n        if col_type != object:\n            c_min = df[col].min()\n            c_max = df[col].max()\n\n            if str(col_type)[:3] == 'int':\n                if c_min > np.iinfo(np.int8).min and c_max < np.iinfo(np.int8).max:\n                    df[col] = df[col].astype(np.int8)\n                elif c_min > np.iinfo(np.int16).min and c_max < np.iinfo(np.int16).max:\n                    df[col] = df[col].astype(np.int16)\n                elif c_min > np.iinfo(np.int32).min and c_max < np.iinfo(np.int32).max:\n                    df[col] = df[col].astype(np.int32)\n            elif str(col_type)[:5] == 'float':\n                # Skip float16 downcasting due to pandas compatibility issues\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.float64)\n\n    end_mem = df.memory_usage().sum() / 1024**2\n    print(f'    Memory usage after optimization: {end_mem:.2f} MB')\n    print(f'    Decreased by {100 * (start_mem - end_mem) / start_mem:.1f}%')\n    return df\n\n\ndef prefilter_transactions(transactions_df, cutoff_date, test_customers, lookback_days=90):\n    \"\"\"Pre-filter transactions to reduce memory footprint (Strategy 2 from design).\"\"\"\n    print(\"\\n[2.5] Pre-filtering transactions to reduce memory...\")\n    original_size = len(transactions_df)\n    original_mem = transactions_df.memory_usage(deep=True).sum() / 1024**2\n    \n    # Filter 1: Keep only lookback window\n    lookback_date = cutoff_date - timedelta(days=lookback_days)\n    filtered_df = transactions_df[\n        (transactions_df['t_dat'] >= lookback_date) &\n        (transactions_df['t_dat'] <= cutoff_date)\n    ].copy()\n    \n    print(f\"    After time filter: {len(filtered_df):,} rows ({100*len(filtered_df)/original_size:.1f}% of original)\")\n    \n    # Filter 2: Keep only test customers\n    test_customer_set = set(test_customers)\n    filtered_df = filtered_df[filtered_df['customer_id'].isin(test_customer_set)].copy()\n    \n    final_size = len(filtered_df)\n    final_mem = filtered_df.memory_usage(deep=True).sum() / 1024**2\n    \n    print(f\"    After customer filter: {final_size:,} rows ({100*final_size/original_size:.1f}% of original)\")\n    print(f\"    Memory reduction: {original_mem:.1f} MB → {final_mem:.1f} MB ({100*(original_mem-final_mem)/original_mem:.1f}% reduction)\")\n    \n    return filtered_df\n\ndef exponential_time_decay(days_ago):\n    \"\"\"Phase 2: Exponential time decay with 14-day half-life.\"\"\"\n    return np.exp(-LAMBDA_DECAY * days_ago)\n\n\ndef get_timescale_multiplier(days_ago):\n    \"\"\"Phase 3: Multi-timescale windows.\"\"\"\n    if days_ago <= 7:\n        return 1.5  # Short-term: trending preferences\n    elif days_ago <= 30:\n        return 1.0  # Medium-term: stable preferences\n    elif days_ago <= 90:\n        return 0.5  # Long-term: general style\n    else:\n        return 0.1  # Very old: minimal weight\n\n\ndef get_recency_multiplier(days_ago):\n    \"\"\"Phase 4: Recency multiplier for purchase frequency boost.\"\"\"\n    if days_ago <= 3:\n        return 2.0  # Last 3 days: 2.0x\n    elif days_ago <= 7:\n        return 1.5  # Last 7 days: 1.5x\n    else:\n        return 1.0\n\n\ndef calculate_repurchase_candidates(transactions_df, cutoff_date, lookback_days=90):\n    \"\"\"Phase 1 Strategy 1: Repurchase candidates from customer history.\"\"\"\n    print(\"  [Strategy 1] Calculating repurchase candidates...\")\n    lookback_date = cutoff_date - timedelta(days=lookback_days)\n    recent_transactions = transactions_df[\n        (transactions_df['t_dat'] >= lookback_date) &\n        (transactions_df['t_dat'] < cutoff_date)\n    ].copy()\n\n    print(f\"    Processing {len(recent_transactions):,} transactions...\")\n\n    if DASK_AVAILABLE and len(recent_transactions) > 1000000 and not DISABLE_DASK:  # Use dask for large datasets\n        print(\"    Using Dask for parallel processing...\")\n\n        # Convert to dask dataframe\n        ddf = dd.from_pandas(recent_transactions, npartitions=mp.cpu_count())\n\n        # Parallel groupby operation\n        result = ddf.groupby(['customer_id', 'article_id']).size().compute()\n\n        # Convert to dictionary format\n        repurchase_candidates = {}\n        for (customer_id, article_id), count in result.items():\n            if customer_id not in repurchase_candidates:\n                repurchase_candidates[customer_id] = {}\n            repurchase_candidates[customer_id][article_id] = count\n\n    else:\n        # Use pandas for smaller datasets or when dask not available\n        print(\"    Using pandas processing...\")\n        repurchase_candidates = {}\n        for customer_id, group in tqdm(recent_transactions.groupby('customer_id'), desc=\"    Processing customers\"):\n            article_counts = group['article_id'].value_counts()\n            repurchase_candidates[customer_id] = article_counts.to_dict()\n\n    print(f\"    Generated repurchase candidates for {len(repurchase_candidates)} customers\")\n    return repurchase_candidates\n\n\ndef calculate_age_group_popular_items(transactions_df, customers_df, cutoff_date, lookback_days=30):\n    \"\"\"Phase 1 Strategy 2: Age group popular items.\"\"\"\n    print(\"  [Strategy 2] Calculating age group popular items...\")\n    lookback_date = cutoff_date - timedelta(days=lookback_days)\n    recent_transactions = transactions_df[\n        (transactions_df['t_dat'] >= lookback_date) &\n        (transactions_df['t_dat'] < cutoff_date)\n    ].copy()\n\n    # Merge with customer age\n    if 'age' not in customers_df.columns:\n        print(\"    Warning: Age column not found, skipping age group strategy\")\n        return {}\n\n    # Create age groups\n    customers_df = customers_df.copy()\n    customers_df['age_group'] = pd.cut(\n        customers_df['age'],\n        bins=[0, 25, 35, 45, 55, 100],\n        labels=['18-25', '26-35', '36-45', '46-55', '55+']\n    )\n\n    print(f\"    Processing {len(recent_transactions):,} transactions for age groups...\")\n\n    if DASK_AVAILABLE and len(recent_transactions) > 500000 and not DISABLE_DASK:  # Use dask for large merges\n        print(\"    Using Dask for parallel merge/groupby...\")\n\n        # Convert to dask dataframes\n        trans_ddf = dd.from_pandas(recent_transactions, npartitions=mp.cpu_count())\n        cust_ddf = dd.from_pandas(customers_df[['customer_id', 'age_group']], npartitions=4)\n\n        # Parallel merge\n        merged_ddf = trans_ddf.merge(cust_ddf, on='customer_id', how='left')\n\n        # Parallel groupby and value_counts equivalent\n        result = merged_ddf.groupby(['age_group', 'article_id']).size().compute()\n\n        # Convert to dictionary format\n        age_group_popular = {}\n        for (age_group, article_id), count in result.items():\n            if pd.notna(age_group):\n                if age_group not in age_group_popular:\n                    age_group_popular[age_group] = {}\n                age_group_popular[age_group][article_id] = count\n\n        # Keep only top 50 per age group\n        for age_group in age_group_popular:\n            sorted_items = sorted(age_group_popular[age_group].items(), key=lambda x: x[1], reverse=True)\n            age_group_popular[age_group] = dict(sorted_items[:50])\n\n    else:\n        # Use pandas for smaller datasets or when dask not available\n        print(\"    Using pandas processing...\")\n        merged = recent_transactions.merge(\n            customers_df[['customer_id', 'age_group']],\n            on='customer_id',\n            how='left'\n        )\n\n        # Calculate popular items per age group\n        age_group_popular = {}\n        for age_group in merged['age_group'].dropna().unique():\n            age_transactions = merged[merged['age_group'] == age_group]\n            popular_items = age_transactions['article_id'].value_counts().head(50).to_dict()\n            age_group_popular[age_group] = popular_items\n\n    print(f\"    Generated popular items for {len(age_group_popular)} age groups\")\n    return age_group_popular\n\n\ndef calculate_cooccurrence_items(transactions_df, cutoff_date, lookback_days=90, min_cooccurrence=5):\n    \"\"\"Phase 1 Strategy 3: Items bought together (co-occurrence) - VECTORIZED VERSION.\"\"\"\n    print(\"  [Strategy 3] Calculating co-occurrence patterns (vectorized)...\")\n    lookback_date = cutoff_date - timedelta(days=lookback_days)\n    recent_transactions = transactions_df[\n        (transactions_df['t_dat'] >= lookback_date) &\n        (transactions_df['t_dat'] < cutoff_date)\n    ].copy()\n\n    print(f\"    Processing {len(recent_transactions):,} transactions for co-occurrence...\")\n\n    # Vectorized approach: Group by customer and date, then create combinations\n    grouped = recent_transactions.groupby(['customer_id', 't_dat'])['article_id'].unique()\n    grouped = grouped.reset_index()\n\n    # Create all combinations for each customer-date group\n    all_combinations = []\n    for articles in tqdm(grouped['article_id'], desc=\"    Creating combinations\"):\n        if len(articles) > 1:\n            # Create unordered pairs (combinations, not permutations)\n            for art1, art2 in combinations(sorted(articles), 2):\n                all_combinations.append((art1, art2))\n\n    if not all_combinations:\n        print(\"    No co-occurrence patterns found\")\n        return {}\n\n    # Convert to DataFrame for efficient counting\n    combo_df = pd.DataFrame(all_combinations, columns=['article1', 'article2'])\n    cooccurrence_counts = combo_df.groupby(['article1', 'article2']).size().reset_index(name='count')\n\n    # Filter by minimum co-occurrence\n    cooccurrence_counts = cooccurrence_counts[cooccurrence_counts['count'] >= min_cooccurrence]\n\n    # Convert to dictionary format\n    cooccurrence = defaultdict(Counter)\n    for _, row in cooccurrence_counts.iterrows():\n        art1, art2, count = row['article1'], row['article2'], row['count']\n        cooccurrence[art1][art2] = count\n        cooccurrence[art2][art1] = count\n\n    # Filter out articles with no co-occurrences\n    filtered_cooccurrence = {k: dict(v) for k, v in cooccurrence.items() if v}\n\n    print(f\"    Generated co-occurrence patterns for {len(filtered_cooccurrence)} articles\")\n    return filtered_cooccurrence\n\n\ndef calculate_category_items(transactions_df, articles_df, cutoff_date, lookback_days=90):\n    \"\"\"Phase 1 Strategy 4: Same category items.\"\"\"\n    print(\"  [Strategy 4] Calculating category-based items...\")\n    lookback_date = cutoff_date - timedelta(days=lookback_days)\n    recent_transactions = transactions_df[\n        (transactions_df['t_dat'] >= lookback_date) &\n        (transactions_df['t_dat'] < cutoff_date)\n    ].copy()\n    \n    # Merge with article categories\n    if 'product_type_no' not in articles_df.columns:\n        print(\"    Warning: product_type_no not found, using article_id grouping\")\n        category_col = 'article_id'\n    else:\n        category_col = 'product_type_no'\n    \n    merged = recent_transactions.merge(\n        articles_df[['article_id', category_col]], \n        on='article_id', \n        how='left'\n    )\n    \n    # Calculate popular items per category\n    category_popular = {}\n    for category, group in merged.groupby(category_col):\n        popular_items = group['article_id'].value_counts().head(50).to_dict()\n        category_popular[category] = popular_items\n    \n    print(f\"    Generated category items for {len(category_popular)} categories\")\n    return category_popular\n\n\ndef calculate_seasonal_trending(transactions_df, cutoff_date, lookback_days=7):\n    \"\"\"Phase 1 Strategy 5: Seasonal trending items.\"\"\"\n    print(\"  [Strategy 5] Calculating seasonal trending items...\")\n    lookback_date = cutoff_date - timedelta(days=lookback_days)\n    recent_transactions = transactions_df[\n        (transactions_df['t_dat'] >= lookback_date) &\n        (transactions_df['t_dat'] < cutoff_date)\n    ].copy()\n    \n    # Calculate trending items (recent popularity surge)\n    trending_items = recent_transactions['article_id'].value_counts().head(100).to_dict()\n    \n    print(f\"    Generated {len(trending_items)} trending items\")\n    return trending_items\n\n\ndef get_product_code_prefix(article_id, articles_df):\n    \"\"\"Get product code prefix for diversity penalty.\"\"\"\n    if 'product_code' in articles_df.columns:\n        product_code = articles_df[articles_df['article_id'] == article_id]['product_code'].values\n        if len(product_code) > 0 and pd.notna(product_code[0]):\n            # Use first 3-4 characters as prefix\n            return str(product_code[0])[:4]\n    return None\n\n\ndef generate_customer_recommendations(\n    customer_id, \n    transactions_df, \n    articles_df, \n    customers_df,\n    cutoff_date,\n    repurchase_candidates,\n    age_group_popular,\n    cooccurrence_items,\n    category_items,\n    seasonal_trending,\n    n=12\n):\n    \"\"\"Generate recommendations for a single customer using all 5 strategies.\"\"\"\n    candidate_scores = defaultdict(float)\n    \n    # Get customer info\n    customer_info = customers_df[customers_df['customer_id'] == customer_id]\n    age_group = None\n    if len(customer_info) > 0 and 'age' in customer_info.columns:\n        age = customer_info['age'].values[0]\n        if pd.notna(age):\n            if age <= 25:\n                age_group = '18-25'\n            elif age <= 35:\n                age_group = '26-35'\n            elif age <= 45:\n                age_group = '36-45'\n            elif age <= 55:\n                age_group = '46-55'\n            else:\n                age_group = '55+'\n    \n    # Strategy 1: Repurchase candidates\n    if customer_id in repurchase_candidates:\n        for article_id, purchase_count in repurchase_candidates[customer_id].items():\n            # Get purchase history for this article\n            customer_history = transactions_df[\n                (transactions_df['customer_id'] == customer_id) &\n                (transactions_df['article_id'] == article_id) &\n                (transactions_df['t_dat'] < cutoff_date)\n            ]\n            \n            if len(customer_history) > 0:\n                # Calculate weighted score\n                total_score = 0.0\n                for _, row in customer_history.iterrows():\n                    days_ago = (cutoff_date - row['t_dat']).days\n                    time_weight = exponential_time_decay(days_ago)\n                    timescale_mult = get_timescale_multiplier(days_ago)\n                    recency_mult = get_recency_multiplier(days_ago)\n                    \n                    score = time_weight * timescale_mult * recency_mult\n                    total_score += score\n                \n                # Phase 4: Purchase frequency boost\n                freq_boost = 1 + 0.2 * purchase_count\n                candidate_scores[article_id] += total_score * freq_boost * 1.0  # Strategy weight\n    \n    # Strategy 2: Age group popular items\n    if age_group and age_group in age_group_popular:\n        for article_id, popularity in age_group_popular[age_group].items():\n            candidate_scores[article_id] += popularity * 0.8  # Strategy weight\n    \n    # Strategy 3: Co-occurrence items\n    customer_articles = transactions_df[\n        (transactions_df['customer_id'] == customer_id) &\n        (transactions_df['t_dat'] < cutoff_date)\n    ]['article_id'].unique()\n    \n    for article_id in customer_articles:\n        if article_id in cooccurrence_items:\n            for cooccur_article, count in cooccurrence_items[article_id].items():\n                candidate_scores[cooccur_article] += count * 1.2  # Strategy weight\n    \n    # Strategy 4: Same category items\n    customer_articles_list = transactions_df[\n        (transactions_df['customer_id'] == customer_id) &\n        (transactions_df['t_dat'] < cutoff_date)\n    ]['article_id'].tolist()\n    \n    if len(customer_articles_list) > 0:\n        # Get categories of customer's articles\n        customer_articles_df = pd.DataFrame({'article_id': customer_articles_list})\n        merged = customer_articles_df.merge(\n            articles_df[['article_id', 'product_type_no']] if 'product_type_no' in articles_df.columns \n            else articles_df[['article_id']],\n            on='article_id',\n            how='left'\n        )\n        \n        if 'product_type_no' in merged.columns:\n            customer_categories = merged['product_type_no'].dropna().unique()\n            for category in customer_categories:\n                if category in category_items:\n                    for article_id, popularity in category_items[category].items():\n                        candidate_scores[article_id] += popularity * 0.6  # Strategy weight\n    \n    # Strategy 5: Seasonal trending\n    for article_id, popularity in seasonal_trending.items():\n        candidate_scores[article_id] += popularity * 0.5  # Strategy weight\n    \n    # Phase 5: Apply diversity penalty\n    if len(candidate_scores) > 0:\n        # Sort by score\n        sorted_candidates = sorted(candidate_scores.items(), key=lambda x: x[1], reverse=True)\n        \n        # Apply diversity penalty\n        final_recommendations = []\n        used_prefixes = set()\n        \n        for article_id, score in sorted_candidates:\n            if len(final_recommendations) >= n:\n                break\n            \n            # Check diversity\n            prefix = get_product_code_prefix(article_id, articles_df)\n            if prefix and prefix in used_prefixes:\n                score *= (1 - DIVERSITY_PENALTY)  # Apply 30% penalty\n            \n            final_recommendations.append((article_id, score))\n            if prefix:\n                used_prefixes.add(prefix)\n        \n        # Sort again after penalty and return top n\n        final_recommendations.sort(key=lambda x: x[1], reverse=True)\n        recommendations = [article_id for article_id, _ in final_recommendations[:n]]\n        \n        # Fill with trending items if needed\n        if len(recommendations) < n:\n            for article_id in seasonal_trending.keys():\n                if article_id not in recommendations:\n                    recommendations.append(article_id)\n                    if len(recommendations) >= n:\n                        break\n        \n        return recommendations[:n]\n    \n    # Fallback: return trending items\n    return list(seasonal_trending.keys())[:n]\n\n\ndef process_customer_batch(args):\n    \"\"\"Worker function for multiprocessing - process a batch of customers.\"\"\"\n    (\n        customer_batch,\n        transactions_df,\n        articles_df,\n        customers_df,\n        cutoff_date,\n        repurchase_candidates,\n        age_group_popular,\n        cooccurrence_items,\n        category_items,\n        seasonal_trending,\n        n\n    ) = args\n\n    batch_recommendations = {}\n    for customer_id in customer_batch:\n        recs = generate_customer_recommendations(\n            customer_id,\n            transactions_df,\n            articles_df,\n            customers_df,\n            cutoff_date,\n            repurchase_candidates,\n            age_group_popular,\n            cooccurrence_items,\n            category_items,\n            seasonal_trending,\n            n=n\n        )\n        batch_recommendations[customer_id] = recs\n\n    return batch_recommendations\n\n\ndef generate_all_recommendations(\n    transactions_df,\n    articles_df,\n    customers_df,\n    customer_ids,\n    cutoff_date,\n    repurchase_candidates,\n    age_group_popular,\n    cooccurrence_items,\n    category_items,\n    seasonal_trending,\n    n=12,\n    batch_size=None,\n    use_multiprocessing=True\n):\n    \"\"\"Generate recommendations for all customers using optimized sequential processing (Strategy 3).\"\"\"\n    if batch_size is None:\n        batch_size = REC_BATCH_SIZE\n    \n    recommendations = {}\n    total_customers = len(customer_ids)\n\n    print(f\"Generating recommendations for {total_customers:,} customers...\")\n\n    # Check memory before starting\n    if PSUTIL_AVAILABLE:\n        mem_gb = get_memory_usage() / 1024\n        print(f\"  Current memory: {mem_gb:.1f} GB\")\n\n    # Use optimized sequential processing (avoids multiprocessing hang)\n    num_batches = (total_customers + batch_size - 1) // batch_size\n    print(f\"  Using optimized sequential processing (batch_size={batch_size})\")\n    print(f\"  Processing {total_customers:,} customers in {num_batches:,} batches...\")\n    print(f\"  Strategy: Avoiding multiprocessing pickle serialization overhead\")\n    \n    # Pre-allocate results dictionary for better memory management\n    recommendations = {cid: [] for cid in customer_ids}\n    \n    # Process in batches with progress tracking\n    start_time = datetime.now()\n    processed_count = 0\n    \n    for batch_idx in tqdm(range(0, total_customers, batch_size), desc=\"  Processing batches\", unit=\"batch\"):\n        batch = customer_ids[batch_idx:batch_idx+batch_size]\n        batch_start = datetime.now()\n        \n        # Process each customer in the batch\n        for customer_id in batch:\n            recs = generate_customer_recommendations(\n                customer_id,\n                transactions_df,\n                articles_df,\n                customers_df,\n                cutoff_date,\n                repurchase_candidates,\n                age_group_popular,\n                cooccurrence_items,\n                category_items,\n                seasonal_trending,\n                n=n\n            )\n            recommendations[customer_id] = recs\n            processed_count += 1\n        \n        # Periodic status updates and memory management\n        batch_num = batch_idx // batch_size\n        if batch_num > 0 and batch_num % 50 == 0:\n            elapsed = (datetime.now() - start_time).total_seconds()\n            rate = processed_count / elapsed\n            remaining = total_customers - processed_count\n            eta_seconds = remaining / rate if rate > 0 else 0\n            eta_minutes = eta_seconds / 60\n            \n            print(f\"\\n    Progress: {processed_count:,}/{total_customers:,} customers ({100*processed_count/total_customers:.1f}%)\")\n            print(f\"    Rate: {rate:.1f} customers/sec | ETA: {eta_minutes:.1f} minutes\")\n            \n            if PSUTIL_AVAILABLE:\n                mem_gb = get_memory_usage() / 1024\n                print(f\"    Memory: {mem_gb:.1f} GB\")\n        \n        # Aggressive garbage collection every 10 batches\n        if batch_num > 0 and batch_num % 10 == 0:\n            gc.collect()\n            \n            # Emergency memory check\n            if PSUTIL_AVAILABLE:\n                mem_gb = get_memory_usage() / 1024\n                if mem_gb > MAX_MEMORY_GB * 0.95:\n                    print(f\"    ⚠️ High memory ({mem_gb:.1f} GB), forcing aggressive GC...\")\n                    gc.collect()\n                    gc.collect()  # Double collection for thorough cleanup\n\n    total_elapsed = (datetime.now() - start_time).total_seconds() / 60\n    print(f\"\\n  ✅ Completed processing {total_customers:,} customers in {total_elapsed:.1f} minutes\")\n    print(f\"  Average rate: {total_customers/total_elapsed:.1f} customers/minute\")\n    \n    return recommendations\n\n\ndef create_submission(recommendations, output_file='submission.csv'):\n    \"\"\"Create submission file in required format.\"\"\"\n    submission_data = []\n    \n    for customer_id, article_list in recommendations.items():\n        prediction = ' '.join([str(article_id).zfill(10) for article_id in article_list])\n        submission_data.append({\n            'customer_id': customer_id,\n            'prediction': prediction\n        })\n    \n    submission_df = pd.DataFrame(submission_data)\n    submission_df.to_csv(output_file, index=False)\n    print(f\"Submission saved to {output_file}\")\n    print(f\"Total predictions: {len(submission_df)}\")\n    \n    return submission_df\n\n\ndef main():\n    \"\"\"Main execution function.\"\"\"\n    start_time = datetime.now()\n    print(\"=\"*80)\n    print(\"H&M Personalized Fashion Recommendations - V2.8 Enhanced\")\n    print(\"=\"*80)\n    print_memory_usage(\"start\")\n\n    # 1. Load Data\n    print(\"\\n[1] Loading data...\")\n    try:\n        # Chunked loading for transactions if file is very large\n        transactions_path = os.path.join(DATA_DIR, \"transactions_train.csv\")\n        if os.path.exists(transactions_path):\n            file_size_gb = os.path.getsize(transactions_path) / (1024**3)\n            print(f\"  Transactions file size: {file_size_gb:.2f} GB\")\n\n            if file_size_gb > 1.0:  # Use chunked loading for large files\n                print(f\"  Using chunked loading for large transactions file (chunk_size={CHUNK_SIZE:,})...\")\n                chunks = []\n                for chunk in pd.read_csv(\n                    transactions_path,\n                    parse_dates=['t_dat'],\n                    chunksize=CHUNK_SIZE\n                ):\n                    chunks.append(chunk)\n                    print(f\"    Loaded chunk with {len(chunk):,} rows\")\n                    if PSUTIL_AVAILABLE:\n                        print_memory_usage(\"after chunk\")\n\n                transactions = pd.concat(chunks, ignore_index=True)\n                del chunks\n                gc.collect()\n            else:\n                transactions = pd.read_csv(transactions_path, parse_dates=['t_dat'])\n        else:\n            raise FileNotFoundError(f\"Transactions file not found: {transactions_path}\")\n\n        articles = pd.read_csv(os.path.join(DATA_DIR, \"articles.csv\"))\n        customers = pd.read_csv(os.path.join(DATA_DIR, \"customers.csv\"))\n        sample_submission = pd.read_csv(os.path.join(DATA_DIR, \"sample_submission.csv\"))\n\n        print(f\"  Transactions: {len(transactions):,} rows\")\n        print(f\"  Articles: {len(articles):,} rows\")\n        print(f\"  Customers: {len(customers):,} rows\")\n        print(f\"  Sample submission: {len(sample_submission):,} rows\")\n        print_memory_usage(\"after loading\")\n\n    except Exception as e:\n        print(f\"  Error loading data: {e}\")\n        sys.exit(1)\n    \n    # Convert IDs to strings\n    transactions['customer_id'] = transactions['customer_id'].astype(str)\n    transactions['article_id'] = transactions['article_id'].astype(str)\n    articles['article_id'] = articles['article_id'].astype(str)\n    customers['customer_id'] = customers['customer_id'].astype(str)\n    sample_submission['customer_id'] = sample_submission['customer_id'].astype(str)\n    \n    # 2. Memory Optimization\n    print(\"\\n[2] Optimizing memory usage...\")\n    transactions = reduce_mem_usage(transactions)\n    articles = reduce_mem_usage(articles)\n    customers = reduce_mem_usage(customers)\n    gc.collect()\n    \n    # 3. Determine cutoff date\n    print(\"\\n[3] Setting up prediction date...\")\n    cutoff_date = transactions['t_dat'].max()\n    test_customers = sample_submission['customer_id'].unique()\n    print(f\"  Cutoff date: {cutoff_date}\")\n    print(f\"  Date range: {transactions['t_dat'].min()} to {cutoff_date}\")\n    print(f\"  Test customers: {len(test_customers):,}\")\n    \n    # 3.5. Pre-filter transactions (CRITICAL FIX: Strategy 2 from design)\n    transactions = prefilter_transactions(\n        transactions, \n        cutoff_date, \n        test_customers, \n        lookback_days=TRANSACTION_LOOKBACK_DAYS\n    )\n    print_memory_usage(\"after pre-filtering\")\n    gc.collect()\n    \n    # 4. Phase 1: Generate candidates using 5 strategies\n    print(\"\\n[4] Phase 1: Enhanced Candidate Generation (5 strategies)...\")\n    print_memory_usage(\"before candidate generation\")\n\n    repurchase_candidates = calculate_repurchase_candidates(transactions, cutoff_date, lookback_days=90)\n    print_memory_usage(\"after repurchase candidates\")\n    if check_memory_limit():\n        print(\"    ⚠️ Continuing despite high memory usage...\")\n    gc.collect()\n\n    age_group_popular = calculate_age_group_popular_items(transactions, customers, cutoff_date, lookback_days=30)\n    print_memory_usage(\"after age group popular\")\n    if check_memory_limit():\n        print(\"    ⚠️ Continuing despite high memory usage...\")\n    gc.collect()\n\n    cooccurrence_items = calculate_cooccurrence_items(transactions, cutoff_date, lookback_days=90, min_cooccurrence=5)\n    print_memory_usage(\"after co-occurrence\")\n    if check_memory_limit():\n        print(\"    ⚠️ Continuing despite high memory usage...\")\n    gc.collect()\n\n    category_items = calculate_category_items(transactions, articles, cutoff_date, lookback_days=90)\n    print_memory_usage(\"after category items\")\n    if check_memory_limit():\n        print(\"    ⚠️ Continuing despite high memory usage...\")\n    gc.collect()\n\n    seasonal_trending = calculate_seasonal_trending(transactions, cutoff_date, lookback_days=7)\n    print_memory_usage(\"after seasonal trending\")\n    if check_memory_limit():\n        print(\"    ⚠️ Continuing despite high memory usage...\")\n    gc.collect()\n\n    # 5. Generate recommendations\n    print(\"\\n[5] Generating recommendations for all customers...\")\n    \n    recommendations = generate_all_recommendations(\n        transactions,\n        articles,\n        customers,\n        test_customers,\n        cutoff_date,\n        repurchase_candidates,\n        age_group_popular,\n        cooccurrence_items,\n        category_items,\n        seasonal_trending,\n        n=12,\n        batch_size=REC_BATCH_SIZE,\n        use_multiprocessing=False  # Explicitly disabled to avoid hang\n    )\n\n    # Memory cleanup - remove large dataframes we don't need anymore\n    del transactions\n    gc.collect()\n    print_memory_usage(\"after cleanup\")\n\n    # 6. Create submission\n    print(\"\\n[6] Creating submission file...\")\n    submission_df = create_submission(recommendations, output_file='submission.csv')\n    \n    # 7. Summary\n    end_time = datetime.now()\n    runtime = (end_time - start_time).total_seconds() / 60\n    final_memory = get_memory_usage()\n\n    print(\"\\n\" + \"=\"*80)\n    print(\"V2.8 ENHANCEMENT COMPLETE - MULTIPROCESSING FIX APPLIED\")\n    print(\"=\"*80)\n    print(f\"Total runtime: {runtime:.2f} minutes\")\n    print(f\"Target met: {'✅' if runtime <= 45 else '⚠️'} (target: <45 min)\")\n    print(f\"Final memory usage: {final_memory:.1f} MB ({final_memory/1024:.1f} GB)\")\n    print(f\"Total predictions: {len(submission_df):,}\")\n    print(f\"Submission file: submission.csv\")\n    print(\"\\nFIXES APPLIED (from design document):\")\n    print(f\"  ✅ Pre-filtered transactions (Strategy 2): 80%+ reduction\")\n    print(f\"  ✅ Optimized sequential processing (Strategy 3): Avoids pickle serialization hang\")\n    print(f\"  ✅ Disabled multiprocessing: Fixes Windows multiprocessing deadlock\")\n    print(f\"  ✅ Vectorized co-occurrence calculation\")\n    print(f\"  ✅ Memory-optimized chunked loading ({CHUNK_SIZE:,} rows)\")\n    print(f\"  ✅ Memory monitoring and limits ({MAX_MEMORY_GB}GB)\")\n    print(f\"  ✅ Aggressive garbage collection and cleanup\")\n    print(f\"  ✅ Dynamic batch sizing (batch_size={REC_BATCH_SIZE})\")\n    if DISABLE_DASK:\n        print(f\"  • Pandas processing (Dask disabled for memory)\")\n    print(\"=\"*80)\n\n    # Final cleanup\n    gc.collect()\n\n    return submission_df\n\n\nif __name__ == \"__main__\":\n    # Required for multiprocessing on Windows\n    mp.freeze_support()\n    submission_df = main()","metadata":{"trusted":true,"execution":{"iopub.status.busy":"2025-11-07T17:27:27.684218Z","iopub.execute_input":"2025-11-07T17:27:27.684521Z","iopub.status.idle":"2025-11-07T17:43:47.903766Z","shell.execute_reply.started":"2025-11-07T17:27:27.684488Z","shell.execute_reply":"2025-11-07T17:43:47.90262Z"}},"outputs":[],"execution_count":null},{"cell_type":"code","source":"\"\"\"\nH&M Personalized Fashion Recommendations - V2.8 Enhanced Implementation\nFollowing the V2 improvement plan for enhanced recommendation system\nExpected MAP@12: 0.032-0.038 (vs 0.012-0.015 baseline)\nRuntime: 27-39 minutes (target: <1 hour)\nOptimized for Kaggle GPU environment\n\"\"\"\n\nimport pandas as pd\nimport numpy as np\nfrom datetime import datetime, timedelta\nfrom collections import defaultdict, Counter\nimport warnings\nimport os\nimport sys\nimport gc\nfrom tqdm import tqdm\nimport multiprocessing as mp\nfrom itertools import combinations\n\n# Try to import optional dependencies\ntry:\n    import psutil\n    PSUTIL_AVAILABLE = True\nexcept ImportError:\n    PSUTIL_AVAILABLE = False\n\n# Try to import dask for parallel dataframe operations\ntry:\n    import dask.dataframe as dd\n    from dask import delayed, compute\n    DASK_AVAILABLE = True\nexcept ImportError:\n    print(\"Warning: Dask not available, using pandas only\")\n    DASK_AVAILABLE = False\n\n# Memory monitoring\ndef get_memory_usage():\n    \"\"\"Get current memory usage in MB.\"\"\"\n    if PSUTIL_AVAILABLE:\n        process = psutil.Process()\n        return process.memory_info().rss / 1024 / 1024\n    else:\n        return 0.0\n\ndef print_memory_usage(message=\"\"):\n    \"\"\"Print current memory usage.\"\"\"\n    if PSUTIL_AVAILABLE:\n        mem_mb = get_memory_usage()\n        mem_gb = mem_mb / 1024\n        warning = \" ⚠️ HIGH MEMORY\" if mem_gb > MAX_MEMORY_GB else \"\"\n        print(f\"    Memory usage{ ' - ' + message if message else ''}: {mem_mb:.1f} MB ({mem_gb:.1f} GB){warning}\")\n\ndef check_memory_limit():\n    \"\"\"Check if memory usage exceeds limit and force garbage collection.\"\"\"\n    if PSUTIL_AVAILABLE:\n        mem_gb = get_memory_usage() / 1024\n        if mem_gb > MAX_MEMORY_GB:\n            print(f\"    ⚠️ Memory usage ({mem_gb:.1f} GB) exceeds limit ({MAX_MEMORY_GB} GB), forcing GC...\")\n            gc.collect()\n            mem_gb_after = get_memory_usage() / 1024\n            print(f\"    Memory after GC: {mem_gb_after:.1f} GB\")\n            return mem_gb_after > MAX_MEMORY_GB\n    return False\n\nwarnings.filterwarnings('ignore')\n\n# Constants\nDATA_DIR = \"/kaggle/input/h-and-m-personalized-fashion-recommendations\"\nLAMBDA_DECAY = np.log(2) / 14  # Exponential decay parameter (14-day half-life)\nDIVERSITY_PENALTY = 0.3  # 30% reduction for same product code prefix\n\n# Memory optimization settings for Kaggle\nCHUNK_SIZE = 500000  # Smaller chunks to reduce memory usage\nMAX_MEMORY_GB = 13  # Target memory limit (increased for Kaggle GPU)\nDISABLE_MULTIPROCESSING = True  # Disable multiprocessing to fix hang issue (Strategy 3)\nDISABLE_DASK = True  # Disable Dask to save memory\nMAX_WORKERS = 2  # Limit workers to avoid memory spikes\nREC_BATCH_SIZE = 1000  # Batch size for recommendations\nTRANSACTION_LOOKBACK_DAYS = 90  # Maximum lookback window for filtering\n\n\ndef reduce_mem_usage(df):\n    \"\"\"Reduce memory usage by downcasting numeric dtypes.\"\"\"\n    start_mem = df.memory_usage().sum() / 1024**2\n\n    for col in df.columns:\n        col_type = df[col].dtype\n\n        # Skip datetime and categorical columns\n        if col_type == 'datetime64[ns]' or col_type.name == 'category':\n            continue\n\n        if col_type != object:\n            c_min = df[col].min()\n            c_max = df[col].max()\n\n            if str(col_type)[:3] == 'int':\n                if c_min > np.iinfo(np.int8).min and c_max < np.iinfo(np.int8).max:\n                    df[col] = df[col].astype(np.int8)\n                elif c_min > np.iinfo(np.int16).min and c_max < np.iinfo(np.int16).max:\n                    df[col] = df[col].astype(np.int16)\n                elif c_min > np.iinfo(np.int32).min and c_max < np.iinfo(np.int32).max:\n                    df[col] = df[col].astype(np.int32)\n            elif str(col_type)[:5] == 'float':\n                # Skip float16 downcasting due to pandas compatibility issues\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.float64)\n\n    end_mem = df.memory_usage().sum() / 1024**2\n    print(f'    Memory usage after optimization: {end_mem:.2f} MB')\n    print(f'    Decreased by {100 * (start_mem - end_mem) / start_mem:.1f}%')\n    return df\n\n\ndef prefilter_transactions(transactions_df, cutoff_date, test_customers, lookback_days=90):\n    \"\"\"Pre-filter transactions to reduce memory footprint (Strategy 2 from design).\"\"\"\n    print(\"\\n[2.5] Pre-filtering transactions to reduce memory...\")\n    original_size = len(transactions_df)\n    original_mem = transactions_df.memory_usage(deep=True).sum() / 1024**2\n    \n    # Filter 1: Keep only lookback window\n    lookback_date = cutoff_date - timedelta(days=lookback_days)\n    filtered_df = transactions_df[\n        (transactions_df['t_dat'] >= lookback_date) &\n        (transactions_df['t_dat'] <= cutoff_date)\n    ].copy()\n    \n    print(f\"    After time filter: {len(filtered_df):,} rows ({100*len(filtered_df)/original_size:.1f}% of original)\")\n    \n    # Filter 2: Keep only test customers\n    test_customer_set = set(test_customers)\n    filtered_df = filtered_df[filtered_df['customer_id'].isin(test_customer_set)].copy()\n    \n    final_size = len(filtered_df)\n    final_mem = filtered_df.memory_usage(deep=True).sum() / 1024**2\n    \n    print(f\"    After customer filter: {final_size:,} rows ({100*final_size/original_size:.1f}% of original)\")\n    print(f\"    Memory reduction: {original_mem:.1f} MB → {final_mem:.1f} MB ({100*(original_mem-final_mem)/original_mem:.1f}% reduction)\")\n    \n    return filtered_df\n\ndef exponential_time_decay(days_ago):\n    \"\"\"Phase 2: Exponential time decay with 14-day half-life.\"\"\"\n    return np.exp(-LAMBDA_DECAY * days_ago)\n\n\ndef get_timescale_multiplier(days_ago):\n    \"\"\"Phase 3: Multi-timescale windows.\"\"\"\n    if days_ago <= 7:\n        return 1.5  # Short-term: trending preferences\n    elif days_ago <= 30:\n        return 1.0  # Medium-term: stable preferences\n    elif days_ago <= 90:\n        return 0.5  # Long-term: general style\n    else:\n        return 0.1  # Very old: minimal weight\n\n\ndef get_recency_multiplier(days_ago):\n    \"\"\"Phase 4: Recency multiplier for purchase frequency boost.\"\"\"\n    if days_ago <= 3:\n        return 2.0  # Last 3 days: 2.0x\n    elif days_ago <= 7:\n        return 1.5  # Last 7 days: 1.5x\n    else:\n        return 1.0\n\n\ndef calculate_repurchase_candidates(transactions_df, cutoff_date, lookback_days=90):\n    \"\"\"Phase 1 Strategy 1: Repurchase candidates from customer history.\"\"\"\n    print(\"  [Strategy 1] Calculating repurchase candidates...\")\n    lookback_date = cutoff_date - timedelta(days=lookback_days)\n    recent_transactions = transactions_df[\n        (transactions_df['t_dat'] >= lookback_date) &\n        (transactions_df['t_dat'] < cutoff_date)\n    ].copy()\n\n    print(f\"    Processing {len(recent_transactions):,} transactions...\")\n\n    if DASK_AVAILABLE and len(recent_transactions) > 1000000 and not DISABLE_DASK:  # Use dask for large datasets\n        print(\"    Using Dask for parallel processing...\")\n\n        # Convert to dask dataframe\n        ddf = dd.from_pandas(recent_transactions, npartitions=mp.cpu_count())\n\n        # Parallel groupby operation\n        result = ddf.groupby(['customer_id', 'article_id']).size().compute()\n\n        # Convert to dictionary format\n        repurchase_candidates = {}\n        for (customer_id, article_id), count in result.items():\n            if customer_id not in repurchase_candidates:\n                repurchase_candidates[customer_id] = {}\n            repurchase_candidates[customer_id][article_id] = count\n\n    else:\n        # Use pandas for smaller datasets or when dask not available\n        print(\"    Using pandas processing...\")\n        repurchase_candidates = {}\n        for customer_id, group in tqdm(recent_transactions.groupby('customer_id'), desc=\"    Processing customers\"):\n            article_counts = group['article_id'].value_counts()\n            repurchase_candidates[customer_id] = article_counts.to_dict()\n\n    print(f\"    Generated repurchase candidates for {len(repurchase_candidates)} customers\")\n    return repurchase_candidates\n\n\ndef calculate_age_group_popular_items(transactions_df, customers_df, cutoff_date, lookback_days=30):\n    \"\"\"Phase 1 Strategy 2: Age group popular items.\"\"\"\n    print(\"  [Strategy 2] Calculating age group popular items...\")\n    lookback_date = cutoff_date - timedelta(days=lookback_days)\n    recent_transactions = transactions_df[\n        (transactions_df['t_dat'] >= lookback_date) &\n        (transactions_df['t_dat'] < cutoff_date)\n    ].copy()\n\n    # Merge with customer age\n    if 'age' not in customers_df.columns:\n        print(\"    Warning: Age column not found, skipping age group strategy\")\n        return {}\n\n    # Create age groups\n    customers_df = customers_df.copy()\n    customers_df['age_group'] = pd.cut(\n        customers_df['age'],\n        bins=[0, 25, 35, 45, 55, 100],\n        labels=['18-25', '26-35', '36-45', '46-55', '55+']\n    )\n\n    print(f\"    Processing {len(recent_transactions):,} transactions for age groups...\")\n\n    if DASK_AVAILABLE and len(recent_transactions) > 500000 and not DISABLE_DASK:  # Use dask for large merges\n        print(\"    Using Dask for parallel merge/groupby...\")\n\n        # Convert to dask dataframes\n        trans_ddf = dd.from_pandas(recent_transactions, npartitions=mp.cpu_count())\n        cust_ddf = dd.from_pandas(customers_df[['customer_id', 'age_group']], npartitions=4)\n\n        # Parallel merge\n        merged_ddf = trans_ddf.merge(cust_ddf, on='customer_id', how='left')\n\n        # Parallel groupby and value_counts equivalent\n        result = merged_ddf.groupby(['age_group', 'article_id']).size().compute()\n\n        # Convert to dictionary format\n        age_group_popular = {}\n        for (age_group, article_id), count in result.items():\n            if pd.notna(age_group):\n                if age_group not in age_group_popular:\n                    age_group_popular[age_group] = {}\n                age_group_popular[age_group][article_id] = count\n\n        # Keep only top 50 per age group\n        for age_group in age_group_popular:\n            sorted_items = sorted(age_group_popular[age_group].items(), key=lambda x: x[1], reverse=True)\n            age_group_popular[age_group] = dict(sorted_items[:50])\n\n    else:\n        # Use pandas for smaller datasets or when dask not available\n        print(\"    Using pandas processing...\")\n        merged = recent_transactions.merge(\n            customers_df[['customer_id', 'age_group']],\n            on='customer_id',\n            how='left'\n        )\n\n        # Calculate popular items per age group\n        age_group_popular = {}\n        for age_group in merged['age_group'].dropna().unique():\n            age_transactions = merged[merged['age_group'] == age_group]\n            popular_items = age_transactions['article_id'].value_counts().head(50).to_dict()\n            age_group_popular[age_group] = popular_items\n\n    print(f\"    Generated popular items for {len(age_group_popular)} age groups\")\n    return age_group_popular\n\n\ndef calculate_cooccurrence_items(transactions_df, cutoff_date, lookback_days=90, min_cooccurrence=5):\n    \"\"\"Phase 1 Strategy 3: Items bought together (co-occurrence) - VECTORIZED VERSION.\"\"\"\n    print(\"  [Strategy 3] Calculating co-occurrence patterns (vectorized)...\")\n    lookback_date = cutoff_date - timedelta(days=lookback_days)\n    recent_transactions = transactions_df[\n        (transactions_df['t_dat'] >= lookback_date) &\n        (transactions_df['t_dat'] < cutoff_date)\n    ].copy()\n\n    print(f\"    Processing {len(recent_transactions):,} transactions for co-occurrence...\")\n\n    # Vectorized approach: Group by customer and date, then create combinations\n    grouped = recent_transactions.groupby(['customer_id', 't_dat'])['article_id'].unique()\n    grouped = grouped.reset_index()\n\n    # Create all combinations for each customer-date group\n    all_combinations = []\n    for articles in tqdm(grouped['article_id'], desc=\"    Creating combinations\"):\n        if len(articles) > 1:\n            # Create unordered pairs (combinations, not permutations)\n            for art1, art2 in combinations(sorted(articles), 2):\n                all_combinations.append((art1, art2))\n\n    if not all_combinations:\n        print(\"    No co-occurrence patterns found\")\n        return {}\n\n    # Convert to DataFrame for efficient counting\n    combo_df = pd.DataFrame(all_combinations, columns=['article1', 'article2'])\n    cooccurrence_counts = combo_df.groupby(['article1', 'article2']).size().reset_index(name='count')\n\n    # Filter by minimum co-occurrence\n    cooccurrence_counts = cooccurrence_counts[cooccurrence_counts['count'] >= min_cooccurrence]\n\n    # Convert to dictionary format\n    cooccurrence = defaultdict(Counter)\n    for _, row in cooccurrence_counts.iterrows():\n        art1, art2, count = row['article1'], row['article2'], row['count']\n        cooccurrence[art1][art2] = count\n        cooccurrence[art2][art1] = count\n\n    # Filter out articles with no co-occurrences\n    filtered_cooccurrence = {k: dict(v) for k, v in cooccurrence.items() if v}\n\n    print(f\"    Generated co-occurrence patterns for {len(filtered_cooccurrence)} articles\")\n    return filtered_cooccurrence\n\n\ndef calculate_category_items(transactions_df, articles_df, cutoff_date, lookback_days=90):\n    \"\"\"Phase 1 Strategy 4: Same category items.\"\"\"\n    print(\"  [Strategy 4] Calculating category-based items...\")\n    lookback_date = cutoff_date - timedelta(days=lookback_days)\n    recent_transactions = transactions_df[\n        (transactions_df['t_dat'] >= lookback_date) &\n        (transactions_df['t_dat'] < cutoff_date)\n    ].copy()\n    \n    # Merge with article categories\n    if 'product_type_no' not in articles_df.columns:\n        print(\"    Warning: product_type_no not found, using article_id grouping\")\n        category_col = 'article_id'\n    else:\n        category_col = 'product_type_no'\n    \n    merged = recent_transactions.merge(\n        articles_df[['article_id', category_col]], \n        on='article_id', \n        how='left'\n    )\n    \n    # Calculate popular items per category\n    category_popular = {}\n    for category, group in merged.groupby(category_col):\n        popular_items = group['article_id'].value_counts().head(50).to_dict()\n        category_popular[category] = popular_items\n    \n    print(f\"    Generated category items for {len(category_popular)} categories\")\n    return category_popular\n\n\ndef calculate_seasonal_trending(transactions_df, cutoff_date, lookback_days=7):\n    \"\"\"Phase 1 Strategy 5: Seasonal trending items.\"\"\"\n    print(\"  [Strategy 5] Calculating seasonal trending items...\")\n    lookback_date = cutoff_date - timedelta(days=lookback_days)\n    recent_transactions = transactions_df[\n        (transactions_df['t_dat'] >= lookback_date) &\n        (transactions_df['t_dat'] < cutoff_date)\n    ].copy()\n    \n    # Calculate trending items (recent popularity surge)\n    trending_items = recent_transactions['article_id'].value_counts().head(100).to_dict()\n    \n    print(f\"    Generated {len(trending_items)} trending items\")\n    return trending_items\n\n\ndef get_product_code_prefix(article_id, articles_df):\n    \"\"\"Get product code prefix for diversity penalty.\"\"\"\n    if 'product_code' in articles_df.columns:\n        product_code = articles_df[articles_df['article_id'] == article_id]['product_code'].values\n        if len(product_code) > 0 and pd.notna(product_code[0]):\n            # Use first 3-4 characters as prefix\n            return str(product_code[0])[:4]\n    return None\n\n\ndef generate_customer_recommendations(\n    customer_id, \n    transactions_df, \n    articles_df, \n    customers_df,\n    cutoff_date,\n    repurchase_candidates,\n    age_group_popular,\n    cooccurrence_items,\n    category_items,\n    seasonal_trending,\n    n=12\n):\n    \"\"\"Generate recommendations for a single customer using all 5 strategies.\"\"\"\n    candidate_scores = defaultdict(float)\n    \n    # Get customer info\n    customer_info = customers_df[customers_df['customer_id'] == customer_id]\n    age_group = None\n    if len(customer_info) > 0 and 'age' in customer_info.columns:\n        age = customer_info['age'].values[0]\n        if pd.notna(age):\n            if age <= 25:\n                age_group = '18-25'\n            elif age <= 35:\n                age_group = '26-35'\n            elif age <= 45:\n                age_group = '36-45'\n            elif age <= 55:\n                age_group = '46-55'\n            else:\n                age_group = '55+'\n    \n    # Strategy 1: Repurchase candidates\n    if customer_id in repurchase_candidates:\n        for article_id, purchase_count in repurchase_candidates[customer_id].items():\n            # Get purchase history for this article\n            customer_history = transactions_df[\n                (transactions_df['customer_id'] == customer_id) &\n                (transactions_df['article_id'] == article_id) &\n                (transactions_df['t_dat'] < cutoff_date)\n            ]\n            \n            if len(customer_history) > 0:\n                # Calculate weighted score\n                total_score = 0.0\n                for _, row in customer_history.iterrows():\n                    days_ago = (cutoff_date - row['t_dat']).days\n                    time_weight = exponential_time_decay(days_ago)\n                    timescale_mult = get_timescale_multiplier(days_ago)\n                    recency_mult = get_recency_multiplier(days_ago)\n                    \n                    score = time_weight * timescale_mult * recency_mult\n                    total_score += score\n                \n                # Phase 4: Purchase frequency boost\n                freq_boost = 1 + 0.2 * purchase_count\n                candidate_scores[article_id] += total_score * freq_boost * 1.0  # Strategy weight\n    \n    # Strategy 2: Age group popular items\n    if age_group and age_group in age_group_popular:\n        for article_id, popularity in age_group_popular[age_group].items():\n            candidate_scores[article_id] += popularity * 0.8  # Strategy weight\n    \n    # Strategy 3: Co-occurrence items\n    customer_articles = transactions_df[\n        (transactions_df['customer_id'] == customer_id) &\n        (transactions_df['t_dat'] < cutoff_date)\n    ]['article_id'].unique()\n    \n    for article_id in customer_articles:\n        if article_id in cooccurrence_items:\n            for cooccur_article, count in cooccurrence_items[article_id].items():\n                candidate_scores[cooccur_article] += count * 1.2  # Strategy weight\n    \n    # Strategy 4: Same category items\n    customer_articles_list = transactions_df[\n        (transactions_df['customer_id'] == customer_id) &\n        (transactions_df['t_dat'] < cutoff_date)\n    ]['article_id'].tolist()\n    \n    if len(customer_articles_list) > 0:\n        # Get categories of customer's articles\n        customer_articles_df = pd.DataFrame({'article_id': customer_articles_list})\n        merged = customer_articles_df.merge(\n            articles_df[['article_id', 'product_type_no']] if 'product_type_no' in articles_df.columns \n            else articles_df[['article_id']],\n            on='article_id',\n            how='left'\n        )\n        \n        if 'product_type_no' in merged.columns:\n            customer_categories = merged['product_type_no'].dropna().unique()\n            for category in customer_categories:\n                if category in category_items:\n                    for article_id, popularity in category_items[category].items():\n                        candidate_scores[article_id] += popularity * 0.6  # Strategy weight\n    \n    # Strategy 5: Seasonal trending\n    for article_id, popularity in seasonal_trending.items():\n        candidate_scores[article_id] += popularity * 0.5  # Strategy weight\n    \n    # Phase 5: Apply diversity penalty\n    if len(candidate_scores) > 0:\n        # Sort by score\n        sorted_candidates = sorted(candidate_scores.items(), key=lambda x: x[1], reverse=True)\n        \n        # Apply diversity penalty\n        final_recommendations = []\n        used_prefixes = set()\n        \n        for article_id, score in sorted_candidates:\n            if len(final_recommendations) >= n:\n                break\n            \n            # Check diversity\n            prefix = get_product_code_prefix(article_id, articles_df)\n            if prefix and prefix in used_prefixes:\n                score *= (1 - DIVERSITY_PENALTY)  # Apply 30% penalty\n            \n            final_recommendations.append((article_id, score))\n            if prefix:\n                used_prefixes.add(prefix)\n        \n        # Sort again after penalty and return top n\n        final_recommendations.sort(key=lambda x: x[1], reverse=True)\n        recommendations = [article_id for article_id, _ in final_recommendations[:n]]\n        \n        # Fill with trending items if needed\n        if len(recommendations) < n:\n            for article_id in seasonal_trending.keys():\n                if article_id not in recommendations:\n                    recommendations.append(article_id)\n                    if len(recommendations) >= n:\n                        break\n        \n        return recommendations[:n]\n    \n    # Fallback: return trending items\n    return list(seasonal_trending.keys())[:n]\n\n\ndef process_customer_batch(args):\n    \"\"\"Worker function for multiprocessing - process a batch of customers.\"\"\"\n    (\n        customer_batch,\n        transactions_df,\n        articles_df,\n        customers_df,\n        cutoff_date,\n        repurchase_candidates,\n        age_group_popular,\n        cooccurrence_items,\n        category_items,\n        seasonal_trending,\n        n\n    ) = args\n\n    batch_recommendations = {}\n    for customer_id in customer_batch:\n        recs = generate_customer_recommendations(\n            customer_id,\n            transactions_df,\n            articles_df,\n            customers_df,\n            cutoff_date,\n            repurchase_candidates,\n            age_group_popular,\n            cooccurrence_items,\n            category_items,\n            seasonal_trending,\n            n=n\n        )\n        batch_recommendations[customer_id] = recs\n\n    return batch_recommendations\n\n\ndef generate_all_recommendations(\n    transactions_df,\n    articles_df,\n    customers_df,\n    customer_ids,\n    cutoff_date,\n    repurchase_candidates,\n    age_group_popular,\n    cooccurrence_items,\n    category_items,\n    seasonal_trending,\n    n=12,\n    batch_size=None,\n    use_multiprocessing=True\n):\n    \"\"\"Generate recommendations for all customers using optimized sequential processing (Strategy 3).\"\"\"\n    if batch_size is None:\n        batch_size = REC_BATCH_SIZE\n    \n    recommendations = {}\n    total_customers = len(customer_ids)\n\n    print(f\"Generating recommendations for {total_customers:,} customers...\")\n\n    # Check memory before starting\n    if PSUTIL_AVAILABLE:\n        mem_gb = get_memory_usage() / 1024\n        print(f\"  Current memory: {mem_gb:.1f} GB\")\n\n    # Use optimized sequential processing (avoids multiprocessing hang)\n    num_batches = (total_customers + batch_size - 1) // batch_size\n    print(f\"  Using optimized sequential processing (batch_size={batch_size})\")\n    print(f\"  Processing {total_customers:,} customers in {num_batches:,} batches...\")\n    print(f\"  Strategy: Avoiding multiprocessing pickle serialization overhead\")\n    \n    # Pre-allocate results dictionary for better memory management\n    recommendations = {cid: [] for cid in customer_ids}\n    \n    # Process in batches with progress tracking\n    start_time = datetime.now()\n    processed_count = 0\n    \n    for batch_idx in tqdm(range(0, total_customers, batch_size), desc=\"  Processing batches\", unit=\"batch\"):\n        batch = customer_ids[batch_idx:batch_idx+batch_size]\n        batch_start = datetime.now()\n        \n        # Process each customer in the batch\n        for customer_id in batch:\n            recs = generate_customer_recommendations(\n                customer_id,\n                transactions_df,\n                articles_df,\n                customers_df,\n                cutoff_date,\n                repurchase_candidates,\n                age_group_popular,\n                cooccurrence_items,\n                category_items,\n                seasonal_trending,\n                n=n\n            )\n            recommendations[customer_id] = recs\n            processed_count += 1\n        \n        # Periodic status updates and memory management\n        batch_num = batch_idx // batch_size\n        if batch_num > 0 and batch_num % 50 == 0:\n            elapsed = (datetime.now() - start_time).total_seconds()\n            rate = processed_count / elapsed\n            remaining = total_customers - processed_count\n            eta_seconds = remaining / rate if rate > 0 else 0\n            eta_minutes = eta_seconds / 60\n            \n            print(f\"\\n    Progress: {processed_count:,}/{total_customers:,} customers ({100*processed_count/total_customers:.1f}%)\")\n            print(f\"    Rate: {rate:.1f} customers/sec | ETA: {eta_minutes:.1f} minutes\")\n            \n            if PSUTIL_AVAILABLE:\n                mem_gb = get_memory_usage() / 1024\n                print(f\"    Memory: {mem_gb:.1f} GB\")\n        \n        # Aggressive garbage collection every 10 batches\n        if batch_num > 0 and batch_num % 10 == 0:\n            gc.collect()\n            \n            # Emergency memory check\n            if PSUTIL_AVAILABLE:\n                mem_gb = get_memory_usage() / 1024\n                if mem_gb > MAX_MEMORY_GB * 0.95:\n                    print(f\"    ⚠️ High memory ({mem_gb:.1f} GB), forcing aggressive GC...\")\n                    gc.collect()\n                    gc.collect()  # Double collection for thorough cleanup\n\n    total_elapsed = (datetime.now() - start_time).total_seconds() / 60\n    print(f\"\\n  ✅ Completed processing {total_customers:,} customers in {total_elapsed:.1f} minutes\")\n    print(f\"  Average rate: {total_customers/total_elapsed:.1f} customers/minute\")\n    \n    return recommendations\n\n\ndef create_submission(recommendations, output_file='submission.csv'):\n    \"\"\"Create submission file in required format.\"\"\"\n    submission_data = []\n    \n    for customer_id, article_list in recommendations.items():\n        prediction = ' '.join([str(article_id).zfill(10) for article_id in article_list])\n        submission_data.append({\n            'customer_id': customer_id,\n            'prediction': prediction\n        })\n    \n    submission_df = pd.DataFrame(submission_data)\n    submission_df.to_csv(output_file, index=False)\n    print(f\"Submission saved to {output_file}\")\n    print(f\"Total predictions: {len(submission_df)}\")\n    \n    return submission_df\n\n\ndef main():\n    \"\"\"Main execution function.\"\"\"\n    start_time = datetime.now()\n    print(\"=\"*80)\n    print(\"H&M Personalized Fashion Recommendations - V2.8 Enhanced\")\n    print(\"=\"*80)\n    print_memory_usage(\"start\")\n\n    # 1. Load Data\n    print(\"\\n[1] Loading data...\")\n    try:\n        # Chunked loading for transactions if file is very large\n        transactions_path = os.path.join(DATA_DIR, \"transactions_train.csv\")\n        if os.path.exists(transactions_path):\n            file_size_gb = os.path.getsize(transactions_path) / (1024**3)\n            print(f\"  Transactions file size: {file_size_gb:.2f} GB\")\n\n            if file_size_gb > 1.0:  # Use chunked loading for large files\n                print(f\"  Using chunked loading for large transactions file (chunk_size={CHUNK_SIZE:,})...\")\n                chunks = []\n                for chunk in pd.read_csv(\n                    transactions_path,\n                    parse_dates=['t_dat'],\n                    chunksize=CHUNK_SIZE\n                ):\n                    chunks.append(chunk)\n                    print(f\"    Loaded chunk with {len(chunk):,} rows\")\n                    if PSUTIL_AVAILABLE:\n                        print_memory_usage(\"after chunk\")\n\n                transactions = pd.concat(chunks, ignore_index=True)\n                del chunks\n                gc.collect()\n            else:\n                transactions = pd.read_csv(transactions_path, parse_dates=['t_dat'])\n        else:\n            raise FileNotFoundError(f\"Transactions file not found: {transactions_path}\")\n\n        articles = pd.read_csv(os.path.join(DATA_DIR, \"articles.csv\"))\n        customers = pd.read_csv(os.path.join(DATA_DIR, \"customers.csv\"))\n        sample_submission = pd.read_csv(os.path.join(DATA_DIR, \"sample_submission.csv\"))\n\n        print(f\"  Transactions: {len(transactions):,} rows\")\n        print(f\"  Articles: {len(articles):,} rows\")\n        print(f\"  Customers: {len(customers):,} rows\")\n        print(f\"  Sample submission: {len(sample_submission):,} rows\")\n        print_memory_usage(\"after loading\")\n\n    except Exception as e:\n        print(f\"  Error loading data: {e}\")\n        sys.exit(1)\n    \n    # Convert IDs to strings\n    transactions['customer_id'] = transactions['customer_id'].astype(str)\n    transactions['article_id'] = transactions['article_id'].astype(str)\n    articles['article_id'] = articles['article_id'].astype(str)\n    customers['customer_id'] = customers['customer_id'].astype(str)\n    sample_submission['customer_id'] = sample_submission['customer_id'].astype(str)\n    \n    # 2. Memory Optimization\n    print(\"\\n[2] Optimizing memory usage...\")\n    transactions = reduce_mem_usage(transactions)\n    articles = reduce_mem_usage(articles)\n    customers = reduce_mem_usage(customers)\n    gc.collect()\n    \n    # 3. Determine cutoff date\n    print(\"\\n[3] Setting up prediction date...\")\n    cutoff_date = transactions['t_dat'].max()\n    test_customers = sample_submission['customer_id'].unique()\n    print(f\"  Cutoff date: {cutoff_date}\")\n    print(f\"  Date range: {transactions['t_dat'].min()} to {cutoff_date}\")\n    print(f\"  Test customers: {len(test_customers):,}\")\n    \n    # 3.5. Pre-filter transactions (CRITICAL FIX: Strategy 2 from design)\n    transactions = prefilter_transactions(\n        transactions, \n        cutoff_date, \n        test_customers, \n        lookback_days=TRANSACTION_LOOKBACK_DAYS\n    )\n    print_memory_usage(\"after pre-filtering\")\n    gc.collect()\n    \n    # 4. Phase 1: Generate candidates using 5 strategies\n    print(\"\\n[4] Phase 1: Enhanced Candidate Generation (5 strategies)...\")\n    print_memory_usage(\"before candidate generation\")\n\n    repurchase_candidates = calculate_repurchase_candidates(transactions, cutoff_date, lookback_days=90)\n    print_memory_usage(\"after repurchase candidates\")\n    if check_memory_limit():\n        print(\"    ⚠️ Continuing despite high memory usage...\")\n    gc.collect()\n\n    age_group_popular = calculate_age_group_popular_items(transactions, customers, cutoff_date, lookback_days=30)\n    print_memory_usage(\"after age group popular\")\n    if check_memory_limit():\n        print(\"    ⚠️ Continuing despite high memory usage...\")\n    gc.collect()\n\n    cooccurrence_items = calculate_cooccurrence_items(transactions, cutoff_date, lookback_days=90, min_cooccurrence=5)\n    print_memory_usage(\"after co-occurrence\")\n    if check_memory_limit():\n        print(\"    ⚠️ Continuing despite high memory usage...\")\n    gc.collect()\n\n    category_items = calculate_category_items(transactions, articles, cutoff_date, lookback_days=90)\n    print_memory_usage(\"after category items\")\n    if check_memory_limit():\n        print(\"    ⚠️ Continuing despite high memory usage...\")\n    gc.collect()\n\n    seasonal_trending = calculate_seasonal_trending(transactions, cutoff_date, lookback_days=7)\n    print_memory_usage(\"after seasonal trending\")\n    if check_memory_limit():\n        print(\"    ⚠️ Continuing despite high memory usage...\")\n    gc.collect()\n\n    # 5. Generate recommendations\n    print(\"\\n[5] Generating recommendations for all customers...\")\n    \n    recommendations = generate_all_recommendations(\n        transactions,\n        articles,\n        customers,\n        test_customers,\n        cutoff_date,\n        repurchase_candidates,\n        age_group_popular,\n        cooccurrence_items,\n        category_items,\n        seasonal_trending,\n        n=12,\n        batch_size=REC_BATCH_SIZE,\n        use_multiprocessing=False  # Explicitly disabled to avoid hang\n    )\n\n    # Memory cleanup - remove large dataframes we don't need anymore\n    del transactions\n    gc.collect()\n    print_memory_usage(\"after cleanup\")\n\n    # 6. Create submission\n    print(\"\\n[6] Creating submission file...\")\n    submission_df = create_submission(recommendations, output_file='submission.csv')\n    \n    # 7. Summary\n    end_time = datetime.now()\n    runtime = (end_time - start_time).total_seconds() / 60\n    final_memory = get_memory_usage()\n\n    print(\"\\n\" + \"=\"*80)\n    print(\"V2.8 ENHANCEMENT COMPLETE - MULTIPROCESSING FIX APPLIED\")\n    print(\"=\"*80)\n    print(f\"Total runtime: {runtime:.2f} minutes\")\n    print(f\"Target met: {'✅' if runtime <= 45 else '⚠️'} (target: <45 min)\")\n    print(f\"Final memory usage: {final_memory:.1f} MB ({final_memory/1024:.1f} GB)\")\n    print(f\"Total predictions: {len(submission_df):,}\")\n    print(f\"Submission file: submission.csv\")\n    print(\"\\nFIXES APPLIED (from design document):\")\n    print(f\"  ✅ Pre-filtered transactions (Strategy 2): 80%+ reduction\")\n    print(f\"  ✅ Optimized sequential processing (Strategy 3): Avoids pickle serialization hang\")\n    print(f\"  ✅ Disabled multiprocessing: Fixes Windows multiprocessing deadlock\")\n    print(f\"  ✅ Vectorized co-occurrence calculation\")\n    print(f\"  ✅ Memory-optimized chunked loading ({CHUNK_SIZE:,} rows)\")\n    print(f\"  ✅ Memory monitoring and limits ({MAX_MEMORY_GB}GB)\")\n    print(f\"  ✅ Aggressive garbage collection and cleanup\")\n    print(f\"  ✅ Dynamic batch sizing (batch_size={REC_BATCH_SIZE})\")\n    if DISABLE_DASK:\n        print(f\"  • Pandas processing (Dask disabled for memory)\")\n    print(\"=\"*80)\n\n    # Final cleanup\n    gc.collect()\n\n    return submission_df\n\n\nif __name__ == \"__main__\":\n    # Required for multiprocessing on Windows\n    mp.freeze_support()\n    submission_df = main()","metadata":{"trusted":true,"execution":{"iopub.status.busy":"2025-11-07T17:43:52.206546Z","iopub.execute_input":"2025-11-07T17:43:52.207325Z"}},"outputs":[],"execution_count":null},{"cell_type":"code","source":"","metadata":{"trusted":true},"outputs":[],"execution_count":null}]}