{"metadata":{"kernelspec":{"language":"python","display_name":"Python 3","name":"python3"},"language_info":{"name":"python","version":"3.11.11","mimetype":"text/x-python","codemirror_mode":{"name":"ipython","version":3},"pygments_lexer":"ipython3","nbconvert_exporter":"python","file_extension":".py"},"kaggle":{"accelerator":"none","dataSources":[{"sourceId":38760,"databundleVersionId":4493939,"sourceType":"competition"}],"dockerImageVersionId":31012,"isInternetEnabled":true,"language":"python","sourceType":"notebook","isGpuEnabled":false}},"nbformat_minor":4,"nbformat":4,"cells":[{"cell_type":"code","source":"# This Python 3 environment comes with many helpful analytics libraries installed\n# It is defined by the kaggle/python Docker image: https://github.com/kaggle/docker-python\n# For example, here's several helpful packages to load\n\nimport numpy as np # linear algebra\nimport pandas as pd # data processing, CSV file I/O (e.g. pd.read_csv)\n\n# Input data files are available in the read-only \"../input/\" directory\n# For example, running this (by clicking run or pressing Shift+Enter) will list all files under the input directory\n\nimport os\nfor dirname, _, filenames in os.walk('/kaggle/input'):\n    for filename in filenames:\n        print(os.path.join(dirname, filename))\n\n# You can write up to 20GB to the current directory (/kaggle/working/) that gets preserved as output when you create a version using \"Save & Run All\" \n# You can also write temporary files to /kaggle/temp/, but they won't be saved outside of the current session","metadata":{"_uuid":"8f2839f25d086af736a60e9eeb907d3b93b6e0e5","_cell_guid":"b1076dfc-b9ad-4769-8c92-a6c4dae69d19","trusted":true,"execution":{"iopub.status.busy":"2025-04-30T11:49:37.697319Z","iopub.execute_input":"2025-04-30T11:49:37.697653Z","iopub.status.idle":"2025-04-30T11:49:38.113944Z","shell.execute_reply.started":"2025-04-30T11:49:37.697619Z","shell.execute_reply":"2025-04-30T11:49:38.113039Z"},"jupyter":{"source_hidden":true}},"outputs":[],"execution_count":null},{"cell_type":"code","source":"import pandas as pd\nimport pyarrow as pa\nimport numpy as np\nimport json\nimport time\nfrom tqdm import tqdm\nfrom collections import Counter, defaultdict\nfrom typing import List, Dict, Set, Tuple, Generator\nimport os\nimport gc\nimport pickle\nfrom IPython.display import display\nimport matplotlib.pyplot as plt\nimport seaborn as sns\nfrom gensim.models import Word2Vec\nimport logging\nimport datetime\nlogging.basicConfig(format='%(asctime)s : %(levelname)s : %(message)s', level=logging.INFO)","metadata":{"trusted":true,"execution":{"iopub.status.busy":"2025-04-30T11:56:48.432140Z","iopub.execute_input":"2025-04-30T11:56:48.432504Z","iopub.status.idle":"2025-04-30T11:56:48.438643Z","shell.execute_reply.started":"2025-04-30T11:56:48.432468Z","shell.execute_reply":"2025-04-30T11:56:48.437803Z"},"jupyter":{"source_hidden":true}},"outputs":[],"execution_count":null},{"cell_type":"code","source":"# ---------------------------------\n# 第1部分: 数据加载和预处理\n# ---------------------------------\n\ndef stream_jsonl(file_path: str, batch_size: int = 5000) -> Generator[List[Dict], None, None]:\n    \"\"\"以批次形式流式读取JSONL文件，避免一次性加载全部内存\"\"\"\n    batch = []\n    with open(file_path, 'r') as f:\n        for line in f:\n            batch.append(json.loads(line))\n            if len(batch) >= batch_size:\n                yield batch\n                batch = []\n        if batch:  # 不要忘记最后一个可能小于batch_size的批次\n            yield batch","metadata":{"trusted":true,"execution":{"iopub.status.busy":"2025-04-30T11:50:16.817281Z","iopub.execute_input":"2025-04-30T11:50:16.817858Z","iopub.status.idle":"2025-04-30T11:50:16.823263Z","shell.execute_reply.started":"2025-04-30T11:50:16.817830Z","shell.execute_reply":"2025-04-30T11:50:16.822413Z"},"jupyter":{"source_hidden":true}},"outputs":[],"execution_count":null},{"cell_type":"code","source":"def stream_events(file_path: str, batch_size: int = 5000) -> Generator[List[Dict], None, None]:\n    \"\"\"从JSONL文件流式读取事件，将每个事件转换为扁平格式\"\"\"\n    for sessions_batch in stream_jsonl(file_path, batch_size):\n        events_batch = []\n        for session in sessions_batch:\n            session_id = session['session']\n            for event in session['events']:\n                events_batch.append({\n                    'session': session_id,\n                    'aid': event['aid'],\n                    'ts': event['ts'],\n                    'type': event['type']\n                })\n        yield events_batch","metadata":{"trusted":true,"execution":{"iopub.status.busy":"2025-04-30T11:50:16.824233Z","iopub.execute_input":"2025-04-30T11:50:16.824574Z","iopub.status.idle":"2025-04-30T11:50:16.859261Z","shell.execute_reply.started":"2025-04-30T11:50:16.824523Z","shell.execute_reply":"2025-04-30T11:50:16.858323Z"},"jupyter":{"source_hidden":true}},"outputs":[],"execution_count":null},{"cell_type":"code","source":"def process_in_batches(data, batch_size=1000, process_func=None, desc=\"处理\"):\n    \"\"\"通用批处理函数，显示批数信息\"\"\"\n    \n    total_items = len(data)\n    total_batches = (total_items + batch_size - 1) // batch_size\n    \n    print(f\"开始{desc}...\")\n    print(f\"总计 {total_items:,} 项，将分为 {total_batches:,} 批 (每批 {batch_size:,} 项)\")\n    \n    results = []\n    start_time = time.time()\n    \n    for batch_idx in range(total_batches):\n        # 计算批次范围\n        start_idx = batch_idx * batch_size\n        end_idx = min(start_idx + batch_size, total_items)\n        \n        # 获取当前批次数据\n        batch_data = data[start_idx:end_idx]\n        \n        # 处理批次\n        print(f\"处理批次 {batch_idx+1}/{total_batches} ({(batch_idx+1)/total_batches*100:.1f}%)...\")\n        if process_func:\n            batch_result = process_func(batch_data)\n            results.append(batch_result)\n        \n        # 显示进度\n        elapsed = time.time() - start_time\n        items_per_second = (batch_idx+1) * batch_size / elapsed if elapsed > 0 else 0\n        remaining_batches = total_batches - (batch_idx+1)\n        remaining_time = (elapsed / (batch_idx+1)) * remaining_batches\n        \n        print(f\"完成进度: {(batch_idx+1)/total_batches*100:.1f}% - \"\n              f\"速度: {items_per_second:.1f} 项/秒 - \"\n              f\"预计剩余时间: {str(datetime.timedelta(seconds=int(remaining_time)))}\")\n    \n    # 处理结果\n    if results and hasattr(results[0], 'append') and callable(getattr(results[0], 'append')):\n        final_result = pd.concat(results, ignore_index=True)\n    else:\n        final_result = results\n    \n    elapsed = time.time() - start_time\n    print(f\"{desc}完成! 总批数: {total_batches}, 总耗时: {elapsed:.1f} 秒\")\n    \n    return final_result","metadata":{"trusted":true,"execution":{"iopub.status.busy":"2025-04-30T11:50:16.860492Z","iopub.execute_input":"2025-04-30T11:50:16.860826Z","iopub.status.idle":"2025-04-30T11:50:16.885681Z","shell.execute_reply.started":"2025-04-30T11:50:16.860792Z","shell.execute_reply":"2025-04-30T11:50:16.884647Z"},"jupyter":{"source_hidden":true}},"outputs":[],"execution_count":null},{"cell_type":"markdown","source":"# ---------------------------------\n# 第2部分: 数据检查与可视化\n# ---------------------------------","metadata":{}},{"cell_type":"code","source":"def peek_jsonl(file_path, num_samples=2):\n    \"\"\"查看JSONL文件的前几个样本\"\"\"\n    samples = []\n    with open(file_path, 'r') as f:\n        for i, line in enumerate(f):\n            if i >= num_samples:\n                break\n            samples.append(json.loads(line))\n    \n    print(f\"=== JSONL文件结构样本 ({file_path}) ===\")\n    for i, sample in enumerate(samples):\n        print(f\"\\n样本 #{i+1}:\")\n        print(json.dumps(sample, indent=2, ensure_ascii=False))\n    \n    # 分析会话结构\n    if samples:\n        sample = samples[0]\n        print(\"\\n=== 会话结构分析 ===\")\n        print(f\"会话ID: {sample.get('session')}\")\n        print(f\"事件数量: {len(sample.get('events', []))}\")\n        \n        if 'events' in sample and sample['events']:\n            event = sample['events'][0]\n            print(f\"\\n事件示例:\")\n            print(f\"  aid (物品ID): {event.get('aid')}\")\n            print(f\"  ts (时间戳): {event.get('ts')}\")\n            print(f\"  type (类型): {event.get('type')}\")","metadata":{"trusted":true,"execution":{"iopub.status.busy":"2025-04-30T11:50:16.889102Z","iopub.execute_input":"2025-04-30T11:50:16.889498Z","iopub.status.idle":"2025-04-30T11:50:16.908138Z","shell.execute_reply.started":"2025-04-30T11:50:16.889460Z","shell.execute_reply":"2025-04-30T11:50:16.907341Z"},"jupyter":{"source_hidden":true}},"outputs":[],"execution_count":null},{"cell_type":"code","source":"def load_and_preview_data(file_path, format='jsonl', batch_size=3, event_limit=10):\n    \"\"\"加载并预览数据，支持JSONL和Parquet格式\"\"\"\n    if format == 'jsonl':\n        # 对于JSONL，加载第一个批次\n        batch = next(stream_jsonl(file_path, batch_size))\n        \n        # 转换为扁平事件格式\n        events = []\n        for session in batch:\n            session_id = session['session']\n            for event in session['events'][:event_limit]:  # 限制每个会话的事件数量\n                events.append({\n                    'session': session_id,\n                    'aid': event['aid'],\n                    'ts': event['ts'],\n                    'type': event['type']\n                })\n        \n        events_df = pd.DataFrame(events)\n        \n        # 显示原始会话数据\n        print(f\"=== 原始会话数据 (前{len(batch)}个会话) ===\")\n        for i, session in enumerate(batch):\n            print(f\"\\n会话 #{i+1}, ID: {session['session']}, 事件数: {len(session['events'])}\")\n            events_sample = session['events'][:event_limit]\n            if len(session['events']) > event_limit:\n                print(f\"(只显示前{event_limit}个事件，共{len(session['events'])}个)\")\n            for j, event in enumerate(events_sample):\n                print(f\"  事件 #{j+1}: 物品={event['aid']}, 类型={event['type']}, 时间戳={event['ts']}\")\n    \n    elif format == 'parquet':\n        # 对于Parquet，直接读取\n        events_df = pd.read_parquet(file_path)\n        # 只取前几行进行展示\n        events_df = events_df.head(batch_size * event_limit)\n    \n    # 显示事件数据框\n    print(f\"\\n=== 扁平化事件数据 ===\")\n    display(events_df.head(10))\n    \n    # 显示数据基本统计信息\n    print(f\"\\n=== 数据统计信息 ===\")\n    print(f\"总事件数: {len(events_df)}\")\n    print(f\"唯一会话数: {events_df['session'].nunique()}\")\n    print(f\"唯一物品数: {events_df['aid'].nunique()}\")\n    print(f\"事件类型分布:\")\n    display(events_df['type'].value_counts())\n    \n    # 可视化事件类型分布\n    plt.figure(figsize=(8, 5))\n    sns.countplot(data=events_df, x='type')\n    plt.title('事件类型分布')\n    plt.xlabel('事件类型')\n    plt.ylabel('数量')\n    plt.xticks(rotation=45)\n    plt.tight_layout()\n    plt.show()\n    \n    return events_df","metadata":{"trusted":true,"execution":{"iopub.status.busy":"2025-04-30T11:50:16.909034Z","iopub.execute_input":"2025-04-30T11:50:16.909334Z","iopub.status.idle":"2025-04-30T11:50:16.933328Z","shell.execute_reply.started":"2025-04-30T11:50:16.909308Z","shell.execute_reply":"2025-04-30T11:50:16.932423Z"},"jupyter":{"source_hidden":true}},"outputs":[],"execution_count":null},{"cell_type":"code","source":"def inspect_features(features_df, feature_type='session', sample_size=5):\n    \"\"\"检查生成的特征\"\"\"\n    print(f\"=== {feature_type.capitalize()}特征检查 ===\")\n    print(f\"特征数量: {len(features_df)}\")\n    print(f\"特征列: {features_df.columns.tolist()}\")\n    print(f\"\\n前{sample_size}个样本:\")\n    display(features_df.head(sample_size))\n    \n    # 显示数值特征的基本统计量\n    numeric_cols = features_df.select_dtypes(include=['number']).columns\n    if len(numeric_cols) > 0:\n        print(f\"\\n数值特征统计:\")\n        display(features_df[numeric_cols].describe())\n        \n        # 可视化特征分布\n        fig, axes = plt.subplots(len(numeric_cols)//3 + 1, 3, figsize=(15, 3*len(numeric_cols)//3 + 3))\n        axes = axes.flatten()\n        \n        for i, col in enumerate(numeric_cols):\n            if i < len(axes):\n                sns.histplot(features_df[col].dropna(), ax=axes[i], kde=True)\n                axes[i].set_title(f'{col} 分布')\n                axes[i].set_xlabel(col)\n        \n        # 隐藏未使用的子图\n        for i in range(len(numeric_cols), len(axes)):\n            axes[i].axis('off')\n        \n        plt.tight_layout()\n        plt.show()","metadata":{"trusted":true,"execution":{"iopub.status.busy":"2025-04-30T11:50:16.934213Z","iopub.execute_input":"2025-04-30T11:50:16.934484Z","iopub.status.idle":"2025-04-30T11:50:16.955258Z","shell.execute_reply.started":"2025-04-30T11:50:16.934460Z","shell.execute_reply":"2025-04-30T11:50:16.954479Z"},"jupyter":{"source_hidden":true}},"outputs":[],"execution_count":null},{"cell_type":"code","source":"def inspect_covisitation_matrix(matrix, top_items=3, top_covisits=5):\n    \"\"\"检查共同访问矩阵\"\"\"\n    print(f\"=== 共同访问矩阵检查 ===\")\n    print(f\"矩阵大小 (唯一物品数): {len(matrix)}\")\n    \n    # 找出共同访问次数最多的几个物品\n    item_weights = [(item_id, sum(counter.values())) for item_id, counter in matrix.items()]\n    top_items_by_weight = sorted(item_weights, key=lambda x: x[1], reverse=True)[:top_items]\n    \n    for item_id, weight in top_items_by_weight:\n        print(f\"\\n物品 {item_id} (总共同访问次数: {weight}):\")\n        top_covisited = matrix[item_id].most_common(top_covisits)\n        for covisit_id, count in top_covisited:\n            print(f\"  - 物品 {covisit_id}: 共同访问 {count} 次\")","metadata":{"trusted":true,"execution":{"iopub.status.busy":"2025-04-30T11:50:16.956222Z","iopub.execute_input":"2025-04-30T11:50:16.956524Z","iopub.status.idle":"2025-04-30T11:50:16.979853Z","shell.execute_reply.started":"2025-04-30T11:50:16.956495Z","shell.execute_reply":"2025-04-30T11:50:16.978633Z"},"jupyter":{"source_hidden":true}},"outputs":[],"execution_count":null},{"cell_type":"code","source":"def inspect_word2vec_model(model, sample_items=3, top_similar=5):\n    \"\"\"检查Word2Vec模型\"\"\"\n    print(f\"=== Word2Vec模型检查 ===\")\n    print(f\"向量维度: {model.vector_size}\")\n    print(f\"词汇量 (唯一物品数): {len(model.wv)}\")\n    \n    # 随机选择几个物品查看其相似物品\n    vocab = list(model.wv.index_to_key)\n    import random\n    sample_vocabs = random.sample(vocab, min(sample_items, len(vocab)))\n    \n    for item_id in sample_vocabs:\n        print(f\"\\n物品 {item_id} 的相似物品:\")\n        try:\n            similar_items = model.wv.most_similar(item_id, topn=top_similar)\n            for similar_id, similarity in similar_items:\n                print(f\"  - 物品 {similar_id}: 相似度 {similarity:.4f}\")\n        except KeyError:\n            print(f\"  物品不在模型词汇表中\")","metadata":{"trusted":true,"execution":{"iopub.status.busy":"2025-04-30T11:50:16.981043Z","iopub.execute_input":"2025-04-30T11:50:16.981543Z","iopub.status.idle":"2025-04-30T11:50:17.002600Z","shell.execute_reply.started":"2025-04-30T11:50:16.981510Z","shell.execute_reply":"2025-04-30T11:50:17.001434Z"},"jupyter":{"source_hidden":true}},"outputs":[],"execution_count":null},{"cell_type":"markdown","source":"# ---------------------------------\n# 第3部分: 特征工程\n# ---------------------------------","metadata":{}},{"cell_type":"code","source":"def create_session_features(events_df: pd.DataFrame) -> pd.DataFrame:\n    \"\"\"为每个会话创建特征\"\"\"\n    # 按会话分组\n    session_features = []\n    \n    for session_id, group in tqdm(events_df.groupby('session')):\n        # 基本会话特征\n        num_events = len(group)\n        session_duration = group['ts'].max() - group['ts'].min() if num_events > 1 else 0\n        unique_items = group['aid'].nunique()\n        \n        # 事件类型统计\n        type_counts = group['type'].value_counts().to_dict()\n        num_clicks = type_counts.get('clicks', 0)\n        num_carts = type_counts.get('carts', 0)\n        num_orders = type_counts.get('orders', 0)\n        \n        # 用户行为特征\n        cart_to_click_ratio = num_carts / num_clicks if num_clicks > 0 else 0\n        order_to_cart_ratio = num_orders / num_carts if num_carts > 0 else 0\n        \n        # 时间相关特征\n        if num_events > 1:\n            timestamps = sorted(group['ts'].values)\n            time_diffs = [timestamps[i+1] - timestamps[i] for i in range(len(timestamps)-1)]\n            avg_time_between_events = np.mean(time_diffs)\n            std_time_between_events = np.std(time_diffs) if len(time_diffs) > 1 else 0\n        else:\n            avg_time_between_events = 0\n            std_time_between_events = 0\n        \n        # 会话最后事件特征\n        sorted_group = group.sort_values('ts')\n        last_event_type = sorted_group.iloc[-1]['type']\n        last_event_ts = sorted_group.iloc[-1]['ts']\n        last_event_aid = sorted_group.iloc[-1]['aid']\n        \n        # 重复商品交互特征\n        repeated_aids = group['aid'].value_counts()\n        max_interactions_for_single_aid = repeated_aids.max() if not repeated_aids.empty else 0\n        num_repeated_aids = sum(repeated_aids > 1)\n        \n        # 将所有特征添加到列表中\n        session_features.append({\n            'session': session_id,\n            'num_events': num_events,\n            'session_duration': session_duration,\n            'unique_items': unique_items,\n            'num_clicks': num_clicks,\n            'num_carts': num_carts,\n            'num_orders': num_orders,\n            'cart_to_click_ratio': cart_to_click_ratio,\n            'order_to_cart_ratio': order_to_cart_ratio,\n            'avg_time_between_events': avg_time_between_events,\n            'std_time_between_events': std_time_between_events,\n            'last_event_type': last_event_type,\n            'last_event_ts': last_event_ts,\n            'last_event_aid': last_event_aid,\n            'max_interactions_for_single_aid': max_interactions_for_single_aid,\n            'num_repeated_aids': num_repeated_aids\n        })\n    \n    return pd.DataFrame(session_features)","metadata":{"trusted":true,"execution":{"iopub.status.busy":"2025-04-30T11:50:17.003940Z","iopub.execute_input":"2025-04-30T11:50:17.004473Z","iopub.status.idle":"2025-04-30T11:50:17.024391Z","shell.execute_reply.started":"2025-04-30T11:50:17.004437Z","shell.execute_reply":"2025-04-30T11:50:17.023346Z"},"jupyter":{"source_hidden":true}},"outputs":[],"execution_count":null},{"cell_type":"code","source":"def create_item_features(events_df: pd.DataFrame) -> pd.DataFrame:\n    \"\"\"为每个商品创建特征\"\"\"\n    # 按商品ID分组\n    item_features = []\n    \n    # 计算全局统计量\n    total_clicks = events_df[events_df['type'] == 'clicks'].shape[0]\n    total_carts = events_df[events_df['type'] == 'carts'].shape[0]\n    total_orders = events_df[events_df['type'] == 'orders'].shape[0]\n    \n    for aid, group in tqdm(events_df.groupby('aid')):\n        # 基本统计量\n        num_sessions = group['session'].nunique()\n        \n        # 事件类型统计\n        type_counts = group['type'].value_counts().to_dict()\n        num_clicks = type_counts.get('clicks', 0)\n        num_carts = type_counts.get('carts', 0)\n        num_orders = type_counts.get('orders', 0)\n        \n        # 转化率特征\n        cart_to_click_ratio = num_carts / num_clicks if num_clicks > 0 else 0\n        order_to_cart_ratio = num_orders / num_carts if num_carts > 0 else 0\n        \n        # 全局流行度特征 (TF-IDF思想)\n        click_popularity = num_clicks / total_clicks if total_clicks > 0 else 0\n        cart_popularity = num_carts / total_carts if total_carts > 0 else 0\n        order_popularity = num_orders / total_orders if total_orders > 0 else 0\n        \n        # 会话覆盖率\n        session_coverage = num_sessions / events_df['session'].nunique()\n        \n        # 时间特征\n        if len(group) > 0:\n            first_ts = group['ts'].min()\n            last_ts = group['ts'].max()\n            avg_ts = group['ts'].mean()\n            ts_span = last_ts - first_ts if last_ts > first_ts else 0\n        else:\n            first_ts = 0\n            last_ts = 0\n            avg_ts = 0\n            ts_span = 0\n        \n        # 将所有特征添加到列表中\n        item_features.append({\n            'aid': aid,\n            'num_sessions': num_sessions,\n            'num_clicks': num_clicks,\n            'num_carts': num_carts,\n            'num_orders': num_orders,\n            'cart_to_click_ratio': cart_to_click_ratio,\n            'order_to_cart_ratio': order_to_cart_ratio,\n            'click_popularity': click_popularity,\n            'cart_popularity': cart_popularity,\n            'order_popularity': order_popularity,\n            'session_coverage': session_coverage,\n            'first_ts': first_ts,\n            'last_ts': last_ts,\n            'avg_ts': avg_ts,\n            'ts_span': ts_span\n        })\n    \n    return pd.DataFrame(item_features)","metadata":{"trusted":true,"execution":{"iopub.status.busy":"2025-04-30T11:50:17.025422Z","iopub.execute_input":"2025-04-30T11:50:17.026220Z","iopub.status.idle":"2025-04-30T11:50:17.047989Z","shell.execute_reply.started":"2025-04-30T11:50:17.026192Z","shell.execute_reply":"2025-04-30T11:50:17.046970Z"},"jupyter":{"source_hidden":true}},"outputs":[],"execution_count":null},{"cell_type":"code","source":"def create_interaction_features(events_df: pd.DataFrame, session_features: pd.DataFrame, item_features: pd.DataFrame) -> pd.DataFrame:\n    \"\"\"创建会话-商品交互特征\"\"\"\n    interaction_features = []\n    \n    # 将会话和商品特征转换为字典，方便查找\n    session_dict = session_features.set_index('session').to_dict('index')\n    item_dict = item_features.set_index('aid').to_dict('index')\n    \n    # 为每个交互创建特征\n    for _, row in tqdm(events_df.iterrows(), total=len(events_df)):\n        session_id = row['session']\n        aid = row['aid']\n        \n        # 基本交互信息\n        interaction = {\n            'session': session_id,\n            'aid': aid,\n            'ts': row['ts'],\n            'type': row['type']\n        }\n        \n        # 添加会话特征\n        if session_id in session_dict:\n            for k, v in session_dict[session_id].items():\n                if k != 'session':  # 跳过ID字段\n                    interaction[f'session_{k}'] = v\n        \n        # 添加商品特征\n        if aid in item_dict:\n            for k, v in item_dict[aid].items():\n                if k != 'aid':  # 跳过ID字段\n                    interaction[f'item_{k}'] = v\n        \n        # 添加交互特征\n        # 例如：该会话中与该商品的交互次数\n        session_item_interactions = events_df[(events_df['session'] == session_id) & (events_df['aid'] == aid)]\n        interaction['interaction_count'] = len(session_item_interactions)\n        \n        # 将特征添加到列表\n        interaction_features.append(interaction)\n    \n    return pd.DataFrame(interaction_features)","metadata":{"trusted":true,"execution":{"iopub.status.busy":"2025-04-30T11:50:17.049022Z","iopub.execute_input":"2025-04-30T11:50:17.049342Z","iopub.status.idle":"2025-04-30T11:50:17.072931Z","shell.execute_reply.started":"2025-04-30T11:50:17.049317Z","shell.execute_reply":"2025-04-30T11:50:17.071940Z"},"jupyter":{"source_hidden":true}},"outputs":[],"execution_count":null},{"cell_type":"markdown","source":"# ---------------------------------\n# 第4部分: 候选生成\n# ---------------------------------","metadata":{}},{"cell_type":"code","source":"def create_covisitation_matrices(events_df: pd.DataFrame, output_dir: str) -> Dict[str, Dict[int, Counter]]:\n    \"\"\"创建不同类型的共同访问矩阵并保存到磁盘\"\"\"\n    os.makedirs(output_dir, exist_ok=True)\n    matrices = {}\n    \n    # 1. 点击到点击的共同访问矩阵（12小时窗口）\n    click_matrix_path = os.path.join(output_dir, 'click_to_click_matrix.pkl')\n    if os.path.exists(click_matrix_path):\n        print(f\"加载现有点击-点击矩阵从 {click_matrix_path}\")\n        with open(click_matrix_path, 'rb') as f:\n            matrices['click_to_click'] = pickle.load(f)\n    else:\n        print(\"创建点击-点击共同访问矩阵...\")\n        click_df = events_df[events_df['type'] == 'clicks'].copy()\n        \n        # 按会话分组并按时间戳排序\n        session_aids = []\n        for session_id, group in click_df.groupby('session'):\n            sorted_group = group.sort_values('ts')\n            session_aids.append((session_id, sorted_group['aid'].tolist(), sorted_group['ts'].tolist()))\n        \n        # 创建共同访问矩阵\n        click_matrix = defaultdict(Counter)\n        time_window = 12 * 60 * 60  # 12小时（秒）\n        \n        for session_id, aids, timestamps in tqdm(session_aids):\n            for i in range(len(aids)):\n                for j in range(len(aids)):\n                    if i != j and abs(timestamps[i] - timestamps[j]) <= time_window:\n                        click_matrix[aids[i]][aids[j]] += 1\n        \n        matrices['click_to_click'] = click_matrix\n        \n        # 保存矩阵\n        with open(click_matrix_path, 'wb') as f:\n            pickle.dump(click_matrix, f)\n    \n    # 2. 点击到购物车/订单的共同访问矩阵（24小时窗口）\n    click_to_cart_matrix_path = os.path.join(output_dir, 'click_to_cart_matrix.pkl')\n    if os.path.exists(click_to_cart_matrix_path):\n        print(f\"加载现有点击-购物车矩阵从 {click_to_cart_matrix_path}\")\n        with open(click_to_cart_matrix_path, 'rb') as f:\n            matrices['click_to_cart'] = pickle.load(f)\n    else:\n        print(\"创建点击-购物车共同访问矩阵...\")\n        # 获取点击事件\n        click_df = events_df[events_df['type'] == 'clicks'].copy()\n        # 获取购物车和订单事件\n        cart_order_df = events_df[events_df['type'].isin(['carts', 'orders'])].copy()\n        \n        # 创建共同访问矩阵\n        click_to_cart_matrix = defaultdict(Counter)\n        time_window = 24 * 60 * 60  # 24小时（秒）\n        \n        # 按会话分组\n        for session_id, group in tqdm(events_df.groupby('session')):\n            # 该会话的点击事件\n            session_clicks = group[group['type'] == 'clicks']\n            # 该会话的购物车/订单事件\n            session_cart_orders = group[group['type'].isin(['carts', 'orders'])]\n            \n            if len(session_clicks) == 0 or len(session_cart_orders) == 0:\n                continue\n            \n            # 对于每个点击事件，查找24小时内的购物车/订单事件\n            for _, click_row in session_clicks.iterrows():\n                click_aid = click_row['aid']\n                click_ts = click_row['ts']\n                \n                for _, cart_order_row in session_cart_orders.iterrows():\n                    cart_order_aid = cart_order_row['aid']\n                    cart_order_ts = cart_order_row['ts']\n                    \n                    # 如果购物车/订单事件发生在点击事件后的24小时内\n                    if 0 <= cart_order_ts - click_ts <= time_window:\n                        click_to_cart_matrix[click_aid][cart_order_aid] += 1\n        \n        matrices['click_to_cart'] = click_to_cart_matrix\n        \n        # 保存矩阵\n        with open(click_to_cart_matrix_path, 'wb') as f:\n            pickle.dump(click_to_cart_matrix, f)\n    \n    return matrices","metadata":{"trusted":true,"execution":{"iopub.status.busy":"2025-04-30T11:50:17.073902Z","iopub.execute_input":"2025-04-30T11:50:17.074173Z","iopub.status.idle":"2025-04-30T11:50:17.094261Z","shell.execute_reply.started":"2025-04-30T11:50:17.074153Z","shell.execute_reply":"2025-04-30T11:50:17.093272Z"},"jupyter":{"source_hidden":true}},"outputs":[],"execution_count":null},{"cell_type":"code","source":"def train_word2vec_embeddings(events_df: pd.DataFrame, output_path: str, vector_size: int = 64) -> Word2Vec:\n    \"\"\"训练Word2Vec嵌入以查找相似商品\"\"\"\n    if os.path.exists(output_path):\n        print(f\"加载现有Word2Vec模型从 {output_path}\")\n        return Word2Vec.load(output_path)\n    \n    print(\"训练Word2Vec模型...\")\n    # 将每个会话转化为商品ID序列\n    sessions = []\n    for session_id, group in tqdm(events_df.groupby('session')):\n        # 按时间戳排序\n        sorted_group = group.sort_values('ts')\n        # 获取序列中的商品ID\n        aid_sequence = sorted_group['aid'].astype(str).tolist()\n        sessions.append(aid_sequence)\n    \n    # 训练Word2Vec模型\n    model = Word2Vec(\n        sentences=sessions,\n        vector_size=vector_size,\n        window=5,\n        min_count=1,\n        workers=4,\n        sg=1,  # 使用Skip-gram\n        epochs=5\n    )\n    \n    # 保存模型\n    model.save(output_path)\n    \n    return model","metadata":{"trusted":true,"execution":{"iopub.status.busy":"2025-04-30T11:50:17.095219Z","iopub.execute_input":"2025-04-30T11:50:17.095545Z","iopub.status.idle":"2025-04-30T11:50:17.115323Z","shell.execute_reply.started":"2025-04-30T11:50:17.095521Z","shell.execute_reply":"2025-04-30T11:50:17.114385Z"},"jupyter":{"source_hidden":true}},"outputs":[],"execution_count":null},{"cell_type":"code","source":"def find_similar_items(model: Word2Vec, item_id: int, top_n: int = 20) -> List[int]:\n    \"\"\"使用Word2Vec模型找到相似商品\"\"\"\n    try:\n        similar_items = model.wv.most_similar(str(item_id), topn=top_n)\n        # 转换回整数并返回\n        return [int(item_id) for item_id, _ in similar_items]\n    except KeyError:\n        # 如果商品ID不在模型中\n        return []","metadata":{"trusted":true,"execution":{"iopub.status.busy":"2025-04-30T11:50:17.116077Z","iopub.execute_input":"2025-04-30T11:50:17.116391Z","iopub.status.idle":"2025-04-30T11:50:17.136579Z","shell.execute_reply.started":"2025-04-30T11:50:17.116361Z","shell.execute_reply":"2025-04-30T11:50:17.135597Z"},"jupyter":{"source_hidden":true}},"outputs":[],"execution_count":null},{"cell_type":"code","source":"def generate_revisit_candidates(session_events: List[Dict]) -> List[int]:\n    \"\"\"根据用户之前的交互生成重访候选集\"\"\"\n    # 获取会话中的所有商品ID\n    aids = [event['aid'] for event in session_events]\n    # 由于用户可能多次与同一商品交互，我们使用集合去重\n    unique_aids = list(set(aids))\n    return unique_aids","metadata":{"trusted":true,"execution":{"iopub.status.busy":"2025-04-30T11:50:17.137661Z","iopub.execute_input":"2025-04-30T11:50:17.137981Z","iopub.status.idle":"2025-04-30T11:50:17.154165Z","shell.execute_reply.started":"2025-04-30T11:50:17.137951Z","shell.execute_reply":"2025-04-30T11:50:17.153213Z"},"jupyter":{"source_hidden":true}},"outputs":[],"execution_count":null},{"cell_type":"code","source":"def generate_candidates_for_session(\n    session_events: List[Dict],\n    covisit_matrices: Dict[str, Dict[int, Counter]],\n    w2v_model: Word2Vec,\n    top_n: int = 20\n) -> Dict[str, List[int]]:\n    \"\"\"为会话生成所有类型的候选项\"\"\"\n    # 获取会话中的商品ID\n    session_aids = [event['aid'] for event in session_events]\n    \n    # 1. 重访候选项\n    revisit_candidates = generate_revisit_candidates(session_events)\n    \n    # 2. 共同访问候选项（点击到点击）\n    click_covisit_candidates = []\n    for aid in session_aids:\n        if aid in covisit_matrices.get('click_to_click', {}):\n            # 获取与该商品最常共同点击的商品\n            top_covisited = [item for item, _ in covisit_matrices['click_to_click'][aid].most_common(top_n)]\n            click_covisit_candidates.extend(top_covisited)\n    \n    # 3. 共同访问候选项（点击到购物车/订单）\n    cart_covisit_candidates = []\n    for aid in session_aids:\n        if aid in covisit_matrices.get('click_to_cart', {}):\n            # 获取用户在点击该商品后最常添加到购物车/下单的商品\n            top_cart_items = [item for item, _ in covisit_matrices['click_to_cart'][aid].most_common(top_n)]\n            cart_covisit_candidates.extend(top_cart_items)\n    \n    # 4. 相似商品候选项（使用Word2Vec）\n    similar_candidates = []\n    for aid in session_aids:\n        similar_items = find_similar_items(w2v_model, aid, top_n=top_n)\n        similar_candidates.extend(similar_items)\n    \n    # 合并所有候选项并按类型分组\n    all_candidates = {\n        'clicks': list(set(revisit_candidates + click_covisit_candidates + similar_candidates)),\n        'carts': list(set(revisit_candidates + click_covisit_candidates + cart_covisit_candidates + similar_candidates)),\n        'orders': list(set(revisit_candidates + cart_covisit_candidates + similar_candidates))\n    }\n    \n    return all_candidates","metadata":{"trusted":true,"execution":{"iopub.status.busy":"2025-04-30T11:50:17.155215Z","iopub.execute_input":"2025-04-30T11:50:17.155605Z","iopub.status.idle":"2025-04-30T11:50:17.172799Z","shell.execute_reply.started":"2025-04-30T11:50:17.155576Z","shell.execute_reply":"2025-04-30T11:50:17.171823Z"},"jupyter":{"source_hidden":true}},"outputs":[],"execution_count":null},{"cell_type":"code","source":"# ---------------------------------\n# 第5部分: 候选排序与特征生成\n# ---------------------------------","metadata":{"trusted":true,"execution":{"iopub.status.busy":"2025-04-30T11:50:17.173841Z","iopub.execute_input":"2025-04-30T11:50:17.174170Z","iopub.status.idle":"2025-04-30T11:50:17.194476Z","shell.execute_reply.started":"2025-04-30T11:50:17.174144Z","shell.execute_reply":"2025-04-30T11:50:17.193612Z"},"jupyter":{"source_hidden":true}},"outputs":[],"execution_count":null},{"cell_type":"code","source":"def create_ranking_features(\n    session_events: List[Dict],\n    candidates: Dict[str, List[int]],\n    covisit_matrices: Dict[str, Dict[int, Counter]],\n    w2v_model: Word2Vec,\n    item_features_df: pd.DataFrame\n) -> Dict[str, pd.DataFrame]:\n    \"\"\"为候选项创建排序特征\"\"\"\n    ranking_features = {}\n    \n    # 获取会话中的商品ID和事件\n    session_aids = [event['aid'] for event in session_events]\n    session_id = session_events[0]['session'] if session_events else None\n    \n    # 按事件类型创建特征\n    for event_type in ['clicks', 'carts', 'orders']:\n        features_list = []\n        \n        for candidate_aid in candidates[event_type]:\n            # 基本特征\n            features = {\n                'session': session_id,\n                'aid': candidate_aid,\n                'event_type': event_type,\n                \n                # 是否在会话中出现过\n                'is_in_session': 1 if candidate_aid in session_aids else 0,\n                \n                # 在会话中出现的次数\n                'occurrence_count': session_aids.count(candidate_aid),\n                \n                # 最近一次出现的位置 (从会话末尾算起)\n                'last_occurrence_pos': len(session_aids) - 1 - session_aids[::-1].index(candidate_aid) if candidate_aid in session_aids else -1,\n                \n                # 共同访问矩阵特征\n                'click_covisit_score': sum(covisit_matrices.get('click_to_click', {}).get(aid, {}).get(candidate_aid, 0) for aid in session_aids),\n                'cart_covisit_score': sum(covisit_matrices.get('click_to_cart', {}).get(aid, {}).get(candidate_aid, 0) for aid in session_aids),\n                \n                # Word2Vec相似度特征\n                'w2v_max_similarity': max([w2v_model.wv.similarity(str(aid), str(candidate_aid)) \n                                          for aid in session_aids \n                                          if str(aid) in w2v_model.wv and str(candidate_aid) in w2v_model.wv] or [0]),\n                'w2v_avg_similarity': np.mean([w2v_model.wv.similarity(str(aid), str(candidate_aid)) \n                                              for aid in session_aids \n                                              if str(aid) in w2v_model.wv and str(candidate_aid) in w2v_model.wv] or [0])\n            }\n            \n            # 添加商品特征\n            if item_features_df is not None and candidate_aid in item_features_df['aid'].values:\n                item_row = item_features_df[item_features_df['aid'] == candidate_aid].iloc[0]\n                for col in item_features_df.columns:\n                    if col != 'aid':\n                        features[f'item_{col}'] = item_row[col]\n            \n            features_list.append(features)\n        \n        # 转换为DataFrame\n        if features_list:\n            ranking_features[event_type] = pd.DataFrame(features_list)\n        else:\n            ranking_features[event_type] = pd.DataFrame()\n    \n    return ranking_features","metadata":{"trusted":true,"execution":{"iopub.status.busy":"2025-04-30T11:50:17.195614Z","iopub.execute_input":"2025-04-30T11:50:17.196442Z","iopub.status.idle":"2025-04-30T11:50:17.220738Z","shell.execute_reply.started":"2025-04-30T11:50:17.196393Z","shell.execute_reply":"2025-04-30T11:50:17.219765Z"},"jupyter":{"source_hidden":true}},"outputs":[],"execution_count":null},{"cell_type":"markdown","source":"# ---------------------------------\n# 第6部分: 模型训练\n# ---------------------------------","metadata":{}},{"cell_type":"code","source":"def train_ranking_model(train_df: pd.DataFrame, features: List[str], target: str = 'is_target'):\n    \"\"\"训练排序模型\"\"\"\n    import lightgbm as lgb\n    \n    # 准备训练数据\n    X = train_df[features]\n    y = train_df[target]\n    \n    # 设置模型参数\n    params = {\n        'objective': 'binary',\n        'metric': 'auc',\n        'boosting_type': 'gbdt',\n        'learning_rate': 0.05,\n        'max_depth': 6,\n        'num_leaves': 31,\n        'feature_fraction': 0.8,\n        'bagging_fraction': 0.8,\n        'bagging_freq': 5,\n        'verbose': -1\n    }\n    \n    # 训练模型\n    print(f\"训练{target}排序模型...\")\n    lgb_train = lgb.Dataset(X, y)\n    model = lgb.train(params, lgb_train, num_boost_round=100)\n    \n    # 特征重要性\n    importance = model.feature_importance(importance_type='gain')\n    feature_imp = pd.DataFrame({'Feature': features, 'Importance': importance})\n    feature_imp = feature_imp.sort_values(by='Importance', ascending=False)\n    \n    print(\"特征重要性:\")\n    display(feature_imp.head(10))\n    \n    # 可视化特征重要性\n    plt.figure(figsize=(10, 6))\n    sns.barplot(x='Importance', y='Feature', data=feature_imp.head(10))\n    plt.title('特征重要性')\n    plt.tight_layout()\n    plt.show()\n    \n    return model","metadata":{"trusted":true,"execution":{"iopub.status.busy":"2025-04-30T11:50:17.224594Z","iopub.execute_input":"2025-04-30T11:50:17.224896Z","iopub.status.idle":"2025-04-30T11:50:17.243887Z","shell.execute_reply.started":"2025-04-30T11:50:17.224875Z","shell.execute_reply":"2025-04-30T11:50:17.242480Z"},"jupyter":{"source_hidden":true}},"outputs":[],"execution_count":null},{"cell_type":"code","source":"def prepare_training_data(events_df: pd.DataFrame, item_features_df: pd.DataFrame, \n                         covisit_matrices: Dict[str, Dict[int, Counter]], w2v_model: Word2Vec,\n                         num_sessions: int = 1000, max_candidates_per_session: int = 100):\n    \"\"\"准备模型训练数据\"\"\"\n    training_data = {\n        'clicks': [],\n        'carts': [],\n        'orders': []\n    }\n    \n    # 获取唯一会话\n    unique_sessions = events_df['session'].unique()\n    if len(unique_sessions) > num_sessions:\n        # 随机选择会话\n        import random\n        selected_sessions = random.sample(list(unique_sessions), num_sessions)\n    else:\n        selected_sessions = unique_sessions\n    \n    # 为每个会话生成训练数据\n    for session_id in tqdm(selected_sessions, desc=\"准备训练数据\"):\n        # 获取该会话的所有事件\n        session_df = events_df[events_df['session'] == session_id].sort_values('ts')\n        \n        # 划分为历史和目标\n        # 使用时间戳的80%作为划分点\n        split_ts = session_df['ts'].min() + 0.8 * (session_df['ts'].max() - session_df['ts'].min())\n        \n        history_df = session_df[session_df['ts'] < split_ts].copy()\n        target_df = session_df[session_df['ts'] >= split_ts].copy()\n        \n        # 如果历史或目标为空，则跳过\n        if len(history_df) == 0 or len(target_df) == 0:\n            continue\n        \n        # 将历史转换为事件列表\n        history_events = [\n            {\n                'session': row['session'],\n                'aid': row['aid'],\n                'ts': row['ts'],\n                'type': row['type']\n            }\n            for _, row in history_df.iterrows()\n        ]\n        \n        # 获取目标商品\n        target_clicks = set(target_df[target_df['type'] == 'clicks']['aid'])\n        target_carts = set(target_df[target_df['type'] == 'carts']['aid'])\n        target_orders = set(target_df[target_df['type'] == 'orders']['aid'])\n        \n        # 生成候选项\n        candidates = generate_candidates_for_session(\n            history_events, covisit_matrices, w2v_model, top_n=max_candidates_per_session\n        )\n        \n        # 为每种事件类型创建排序特征\n        ranking_features = create_ranking_features(\n            history_events, candidates, covisit_matrices, w2v_model, item_features_df\n        )\n        \n        # 添加目标标签并加入训练数据\n        for event_type in ['clicks', 'carts', 'orders']:\n            if event_type in ranking_features and not ranking_features[event_type].empty:\n                df = ranking_features[event_type].copy()\n                \n                # 添加目标标签\n                if event_type == 'clicks':\n                    df['is_target'] = df['aid'].isin(target_clicks).astype(int)\n                elif event_type == 'carts':\n                    df['is_target'] = df['aid'].isin(target_carts).astype(int)\n                else:  # orders\n                    df['is_target'] = df['aid'].isin(target_orders).astype(int)\n                \n                # 添加到训练数据\n                training_data[event_type].append(df)\n    \n    # 合并所有会话的训练数据\n    for event_type in training_data:\n        if training_data[event_type]:\n            training_data[event_type] = pd.concat(training_data[event_type], ignore_index=True)\n        else:\n            training_data[event_type] = pd.DataFrame()\n        \n        print(f\"{event_type}事件训练数据大小: {len(training_data[event_type])}\")\n        print(f\"正样本比例: {training_data[event_type]['is_target'].mean():.4f}\")\n    \n    return training_data","metadata":{"trusted":true,"execution":{"iopub.status.busy":"2025-04-30T11:50:17.244949Z","iopub.execute_input":"2025-04-30T11:50:17.245257Z","iopub.status.idle":"2025-04-30T11:50:17.269798Z","shell.execute_reply.started":"2025-04-30T11:50:17.245225Z","shell.execute_reply":"2025-04-30T11:50:17.268763Z"},"jupyter":{"source_hidden":true}},"outputs":[],"execution_count":null},{"cell_type":"markdown","source":"# ---------------------------------\n# 第7部分: 预测和提交\n# ---------------------------------","metadata":{}},{"cell_type":"code","source":"def predict_for_session(\n    session_events: List[Dict],\n    covisit_matrices: Dict[str, Dict[int, Counter]],\n    w2v_model: Word2Vec,\n    item_features_df: pd.DataFrame,\n    models: Dict[str, object],\n    top_n: int = 20\n) -> Dict[str, List[int]]:\n    \"\"\"为会话生成预测\"\"\"\n    # 生成候选项\n    candidates = generate_candidates_for_session(\n        session_events, covisit_matrices, w2v_model, top_n=100  # 生成更多候选项以便排序\n    )\n    \n    # 为候选项创建排序特征\n    ranking_features = create_ranking_features(\n        session_events, candidates, covisit_matrices, w2v_model, item_features_df\n    )\n    \n    # 使用模型预测各个候选项的分数\n    predictions = {}\n    \n    for event_type in ['clicks', 'carts', 'orders']:\n        if event_type in ranking_features and not ranking_features[event_type].empty:\n            df = ranking_features[event_type]\n            \n            # 提取特征列\n            feature_cols = [col for col in df.columns if col not in ['session', 'aid', 'event_type']]\n            \n            # 预测分数\n            if event_type in models and models[event_type] is not None:\n                scores = models[event_type].predict(df[feature_cols])\n                \n                # 将预测分数与商品ID配对并排序\n                scored_items = list(zip(df['aid'].values, scores))\n                sorted_items = sorted(scored_items, key=lambda x: x[1], reverse=True)\n                \n                # 获取前top_n个商品\n                predictions[event_type] = [item[0] for item in sorted_items[:top_n]]\n            else:\n                # 如果没有模型，则使用基于规则的候选项\n                predictions[event_type] = candidates[event_type][:top_n]\n        else:\n            # 如果没有候选项，则返回空列表\n            predictions[event_type] = []\n    \n    return predictions","metadata":{"trusted":true,"execution":{"iopub.status.busy":"2025-04-30T11:50:17.270996Z","iopub.execute_input":"2025-04-30T11:50:17.271344Z","iopub.status.idle":"2025-04-30T11:50:17.298525Z","shell.execute_reply.started":"2025-04-30T11:50:17.271319Z","shell.execute_reply":"2025-04-30T11:50:17.297524Z"},"jupyter":{"source_hidden":true}},"outputs":[],"execution_count":null},{"cell_type":"code","source":"def create_submission(\n    test_path: str,\n    covisit_matrices: Dict[str, Dict[int, Counter]],\n    w2v_model: Word2Vec,\n    item_features_df: pd.DataFrame,\n    models: Dict[str, object],\n    output_path: str = 'submission.csv',\n    batch_size: int = 100\n) -> None:\n    \"\"\"创建比赛提交文件\"\"\"\n    # 打开输出文件\n    with open(output_path, 'w') as f:\n        # 写入标题行\n        f.write('session_type,labels\\n')\n        \n        # 分批处理测试数据\n        batch_count = 0\n        for session_batch in stream_jsonl(test_path, batch_size=batch_size):\n            batch_count += 1\n            print(f\"处理测试批次 {batch_count}...\")\n            \n            for session in tqdm(session_batch, desc=f\"批次 {batch_count}\"):\n                session_id = session['session']\n                \n                # 预测\n                predictions = predict_for_session(\n                    session['events'], covisit_matrices, w2v_model, \n                    item_features_df, models\n                )\n                \n                # 写入预测结果\n                for event_type in ['clicks', 'carts', 'orders']:\n                    # 确保至少有一些预测\n                    preds = predictions.get(event_type, [])\n                    if len(preds) < 20:\n                        # 如果预测不足20个，则使用基于规则的方法补充\n                        rule_based = generate_candidates_for_session(\n                            session['events'], covisit_matrices, w2v_model\n                        )[event_type]\n                        \n                        # 添加不在当前预测中的商品\n                        for item in rule_based:\n                            if item not in preds and len(preds) < 20:\n                                preds.append(item)\n                    \n                    # 格式化预测结果\n                    labels_str = ' '.join(map(str, preds[:20]))  # 最多20个预测\n                    f.write(f\"{session_id};{event_type},{labels_str}\\n\")\n            \n            # 释放内存\n            gc.collect()\n    \n    print(f\"提交文件已保存至 {output_path}\")","metadata":{"trusted":true,"execution":{"iopub.status.busy":"2025-04-30T11:50:17.299518Z","iopub.execute_input":"2025-04-30T11:50:17.299849Z","iopub.status.idle":"2025-04-30T11:50:17.323147Z","shell.execute_reply.started":"2025-04-30T11:50:17.299826Z","shell.execute_reply":"2025-04-30T11:50:17.322242Z"},"jupyter":{"source_hidden":true}},"outputs":[],"execution_count":null},{"cell_type":"code","source":"def convert_jsonl_to_parquet(jsonl_path: str, parquet_path: str, batch_size: int = 10000) -> None:\n    \"\"\"Convert a JSONL file to Parquet format for better I/O performance.\"\"\"\n    import os\n    import pandas as pd\n    \n    print(f\"Converting {jsonl_path} to Parquet format...\")\n    \n    # Create directory if it doesn't exist\n    os.makedirs(os.path.dirname(parquet_path), exist_ok=True)\n    \n    # Remove the file if it exists\n    if os.path.exists(parquet_path):\n        print(f\"Removing existing file: {parquet_path}\")\n        os.remove(parquet_path)\n    \n    # Process in batches\n    all_data = []\n    batch_idx = 0\n    \n    for events_batch in stream_events(jsonl_path, batch_size=batch_size):\n        all_data.extend(events_batch)\n        batch_idx += 1\n        print(f\"Processed batch {batch_idx}\")\n    \n    # Convert to DataFrame and save as parquet\n    df = pd.DataFrame(all_data)\n    df.to_parquet(\n        parquet_path, \n        engine='pyarrow', \n        index=False, \n        compression='snappy'\n    )\n    \n    print(f\"Conversion completed: {parquet_path}\")","metadata":{"trusted":true,"execution":{"iopub.status.busy":"2025-04-30T12:04:45.421226Z","iopub.execute_input":"2025-04-30T12:04:45.421592Z","iopub.status.idle":"2025-04-30T12:04:45.428418Z","shell.execute_reply.started":"2025-04-30T12:04:45.421563Z","shell.execute_reply":"2025-04-30T12:04:45.427394Z"}},"outputs":[],"execution_count":null},{"cell_type":"markdown","source":"# ---------------------------------\n# 第8部分: 评估功能\n# ---------------------------------","metadata":{}},{"cell_type":"code","source":"def evaluate_predictions(\n    truth_df: pd.DataFrame,\n    pred_df: pd.DataFrame,\n    weight_clicks: float = 0.10,\n    weight_carts: float = 0.30,\n    weight_orders: float = 0.60\n) -> Dict[str, float]:\n    \"\"\"评估预测结果\"\"\"\n    results = {}\n    \n    # 为每种类型计算召回率\n    for event_type in ['clicks', 'carts', 'orders']:\n        type_truth = truth_df[truth_df['type'] == event_type]\n        \n        # 会话和真实标签的映射\n        session_labels = {}\n        for _, row in type_truth.iterrows():\n            session_id = row['session']\n            aid = row['aid']\n            if session_id not in session_labels:\n                session_labels[session_id] = []\n            session_labels[session_id].append(aid)\n        \n        # 获取预测\n        type_preds = pred_df[pred_df['type'] == event_type]\n        \n        # 计算召回率\n        recall_numerator = 0\n        recall_denominator = 0\n        \n        for session_id, true_aids in session_labels.items():\n            # 获取该会话的预测\n            session_preds = type_preds[type_preds['session'] == session_id]\n            if len(session_preds) > 0:\n                pred_aids = session_preds.iloc[0]['preds']\n                # 计算交集大小\n                hits = len(set(true_aids).intersection(set(pred_aids)))\n                recall_numerator += hits\n                recall_denominator += min(20, len(true_aids))\n        \n        # 计算最终召回率\n        if recall_denominator > 0:\n            recall = recall_numerator / recall_denominator\n        else:\n            recall = 0\n        \n        results[f'recall_{event_type}'] = recall\n    \n    # 计算加权得分\n    weighted_score = (\n        results.get('recall_clicks', 0) * weight_clicks +\n        results.get('recall_carts', 0) * weight_carts +\n        results.get('recall_orders', 0) * weight_orders\n    )\n    results['weighted_score'] = weighted_score\n    \n    return results    ","metadata":{"trusted":true,"execution":{"iopub.status.busy":"2025-04-30T08:17:45.707102Z","iopub.execute_input":"2025-04-30T08:17:45.707437Z","iopub.status.idle":"2025-04-30T08:17:45.818143Z","shell.execute_reply.started":"2025-04-30T08:17:45.707414Z","shell.execute_reply":"2025-04-30T08:17:45.817097Z"},"jupyter":{"source_hidden":true}},"outputs":[],"execution_count":null},{"cell_type":"code","source":"\"\"\"主函数\"\"\"\n# 路径配置\nworking_dir = '/kaggle/working/'\ndata_dir = '/kaggle/input/otto-recommender-system/'\nprocessed_dir = os.path.join(working_dir, 'processed')\nmatrices_dir = os.path.join(processed_dir, 'matrices')\nmodels_dir = os.path.join(processed_dir, 'models')\nos.makedirs(processed_dir, exist_ok=True)\nos.makedirs(matrices_dir, exist_ok=True)\nos.makedirs(models_dir, exist_ok=True)\n    \ntrain_path = os.path.join(data_dir, 'train.jsonl')\ntest_path = os.path.join(data_dir, 'test.jsonl')\n    \n# 1. 数据预处理\n# 如果需要，将JSONL转换为Parquet\ntrain_parquet = os.path.join(processed_dir, 'train_events.parquet')\ntest_parquet = os.path.join(processed_dir, 'test_events.parquet')","metadata":{"trusted":true,"execution":{"iopub.status.busy":"2025-04-30T11:50:49.306410Z","iopub.execute_input":"2025-04-30T11:50:49.307204Z","iopub.status.idle":"2025-04-30T11:50:49.313997Z","shell.execute_reply.started":"2025-04-30T11:50:49.307166Z","shell.execute_reply":"2025-04-30T11:50:49.312955Z"},"jupyter":{"source_hidden":true}},"outputs":[],"execution_count":null},{"cell_type":"code","source":"# if not os.path.exists(train_parquet):\nconvert_jsonl_to_parquet(train_path, train_parquet, batch_size=5000)","metadata":{"trusted":true,"execution":{"iopub.status.busy":"2025-04-30T12:04:49.786181Z","iopub.execute_input":"2025-04-30T12:04:49.786476Z"},"jupyter":{"outputs_hidden":true},"collapsed":true},"outputs":[],"execution_count":null},{"cell_type":"code","source":"if not os.path.exists(test_parquet):\n    convert_jsonl_to_parquet(test_path, test_parquet, batch_size=5000)","metadata":{"trusted":true,"execution":{"iopub.status.busy":"2025-04-30T11:59:19.662309Z","iopub.execute_input":"2025-04-30T11:59:19.662699Z","iopub.status.idle":"2025-04-30T11:59:19.667454Z","shell.execute_reply.started":"2025-04-30T11:59:19.662669Z","shell.execute_reply":"2025-04-30T11:59:19.666616Z"}},"outputs":[],"execution_count":null},{"cell_type":"code","source":"\n    # # 2. 加载数据\n    # print(\"加载训练数据...\")\n    # train_df = pd.read_parquet(train_parquet)\n    \n    # # 对于大数据集，可以取样本加速开发过程\n    # # train_df = train_df.sample(frac=0.1, random_state=42)\n    \n    # # 3. 特征工程\n    # # 创建会话特征\n    # session_features_path = os.path.join(processed_dir, 'session_features.parquet')\n    # if os.path.exists(session_features_path):\n    #     session_features = pd.read_parquet(session_features_path)\n    # else:\n    #     print(\"创建会话特征...\")\n    #     session_features = create_session_features(train_df)\n    #     session_features.to_parquet(session_features_path, index=False)\n    \n    # # 创建商品特征\n    # item_features_path = os.path.join(processed_dir, 'item_features.parquet')\n    # if os.path.exists(item_features_path):\n    #     item_features = pd.read_parquet(item_features_path)\n    # else:\n    #     print(\"创建商品特征...\")\n    #     item_features = create_item_features(train_df)\n    #     item_features.to_parquet(item_features_path, index=False)\n    \n    # # 4. 候选生成\n    # # 创建共同访问矩阵\n    # print(\"创建共同访问矩阵...\")\n    # covisit_matrices = create_covisitation_matrices(train_df, matrices_dir)\n    \n    # # 训练Word2Vec模型\n    # w2v_path = os.path.join(processed_dir, 'word2vec_model.bin')\n    # w2v_model = train_word2vec_embeddings(train_df, w2v_path, vector_size=64)\n    \n    # # 5. 准备训练数据\n    # # 这一步可能会消耗大量内存，根据需要调整参数\n    # training_data_path = os.path.join(processed_dir, 'training_data.pkl')\n    # if os.path.exists(training_data_path):\n    #     print(f\"加载现有训练数据从 {training_data_path}\")\n    #     with open(training_data_path, 'rb') as f:\n    #         training_data = pickle.load(f)\n    # else:\n    #     print(\"准备训练数据...\")\n    #     training_data = prepare_training_data(\n    #         train_df, item_features, covisit_matrices, w2v_model,\n    #         num_sessions=10000, max_candidates_per_session=100\n    #     )\n        \n    #     # 保存训练数据\n    #     with open(training_data_path, 'wb') as f:\n    #         pickle.dump(training_data, f)\n    \n    # # 6. 训练排序模型\n    # models = {}\n    \n    # for event_type in ['clicks', 'carts', 'orders']:\n    #     # 如果有足够的训练数据，则训练模型\n    #     if event_type in training_data and len(training_data[event_type]) > 0:\n    #         df = training_data[event_type]\n            \n    #         # 选择特征列\n    #         feature_cols = [col for col in df.columns \n    #                         if col not in ['session', 'aid', 'event_type', 'is_target']]\n            \n    #         # 训练模型\n    #         model_path = os.path.join(models_dir, f'{event_type}_model.pkl')\n    #         if os.path.exists(model_path):\n    #             print(f\"加载现有{event_type}模型从 {model_path}\")\n    #             with open(model_path, 'rb') as f:\n    #                 models[event_type] = pickle.load(f)\n    #         else:\n    #             print(f\"训练{event_type}模型...\")\n    #             models[event_type] = train_ranking_model(df, feature_cols, 'is_target')\n                \n    #             # 保存模型\n    #             with open(model_path, 'wb') as f:\n    #                 pickle.dump(models[event_type], f)\n    #     else:\n    #         print(f\"没有足够的{event_type}训练数据，将使用基于规则的方法\")\n    #         models[event_type] = None\n    \n    # # 7. 创建提交文件\n    # submission_path = 'submission.csv'\n    # print(\"创建提交文件...\")\n    # create_submission(\n    #     test_path, covisit_matrices, w2v_model, item_features, models,\n    #     output_path=submission_path, batch_size=100\n    # )\n    \n    # print(f\"完成！提交文件已保存至 {submission_path}\")","metadata":{"trusted":true,"jupyter":{"source_hidden":true}},"outputs":[],"execution_count":null}]}