{"metadata":{"kernelspec":{"language":"python","display_name":"Python 3","name":"python3"},"language_info":{"name":"python","version":"3.10.14","mimetype":"text/x-python","codemirror_mode":{"name":"ipython","version":3},"pygments_lexer":"ipython3","nbconvert_exporter":"python","file_extension":".py"},"kaggle":{"accelerator":"gpu","dataSources":[{"sourceId":84493,"databundleVersionId":9871156,"sourceType":"competition"}],"dockerImageVersionId":30786,"isInternetEnabled":false,"language":"python","sourceType":"notebook","isGpuEnabled":true}},"nbformat_minor":4,"nbformat":4,"cells":[{"cell_type":"code","source":"# This Python 3 environment comes with many helpful analytics libraries installed\n# It is defined by the kaggle/python Docker image: https://github.com/kaggle/docker-python\n# For example, here's several helpful packages to load\n\nimport numpy as np # linear algebra\nimport pandas as pd # data processing, CSV file I/O (e.g. pd.read_csv)\n\n# Input data files are available in the read-only \"../input/\" directory\n# For example, running this (by clicking run or pressing Shift+Enter) will list all files under the input directory\n\nimport os\nfor dirname, _, filenames in os.walk('/kaggle/input'):\n    for filename in filenames:\n        print(os.path.join(dirname, filename))\n\n# You can write up to 20GB to the current directory (/kaggle/working/) that gets preserved as output when you create a version using \"Save & Run All\" \n# You can also write temporary files to /kaggle/temp/, but they won't be saved outside of the current session","metadata":{"_uuid":"8f2839f25d086af736a60e9eeb907d3b93b6e0e5","_cell_guid":"b1076dfc-b9ad-4769-8c92-a6c4dae69d19","execution":{"iopub.status.busy":"2024-10-27T08:53:43.822279Z","iopub.execute_input":"2024-10-27T08:53:43.822721Z","iopub.status.idle":"2024-10-27T08:53:44.953959Z","shell.execute_reply.started":"2024-10-27T08:53:43.822681Z","shell.execute_reply":"2024-10-27T08:53:44.952511Z"},"trusted":true},"execution_count":null,"outputs":[]},{"cell_type":"code","source":"!nvidia-smi","metadata":{"execution":{"iopub.status.busy":"2024-10-27T08:53:44.956471Z","iopub.execute_input":"2024-10-27T08:53:44.957166Z","iopub.status.idle":"2024-10-27T08:53:46.155353Z","shell.execute_reply.started":"2024-10-27T08:53:44.957105Z","shell.execute_reply":"2024-10-27T08:53:46.153835Z"},"trusted":true},"execution_count":null,"outputs":[]},{"cell_type":"code","source":"import os\nimport sys\nimport pandas as pd\nimport polars as pl\n\nsys.path.append('/kaggle/input/jane-street-real-time-market-data-forecasting/')\nfrom kaggle_evaluation import jane_street_inference_server","metadata":{"execution":{"iopub.status.busy":"2024-10-27T08:53:56.006885Z","iopub.execute_input":"2024-10-27T08:53:56.007366Z","iopub.status.idle":"2024-10-27T08:53:56.472145Z","shell.execute_reply.started":"2024-10-27T08:53:56.007312Z","shell.execute_reply":"2024-10-27T08:53:56.471170Z"},"trusted":true},"execution_count":null,"outputs":[]},{"cell_type":"code","source":"import joblib\nimport polars as pl\nimport pandas as pd\nimport numpy as np\nfrom sklearn.preprocessing import StandardScaler\n\ndef load_model(symbol_id: int):\n    \"\"\"加载本地保存的 LightGBM 模型.\"\"\"\n    model_path = f\"model_{symbol_id}.joblib\"\n    if os.path.exists(model_path):\n        return joblib.load(model_path)\n    else:\n        print(f\"No model found for symbol {symbol_id}\")\n        return None\n\ndef predict(test: pl.DataFrame, lags: pl.DataFrame | None = None) -> pl.DataFrame:\n    \"\"\"Make a prediction using LightGBM models.\"\"\"\n    global feature_columns\n\n    # 转换 Polars DataFrame 为 Pandas DataFrame\n    test_df = test.to_pandas()\n    predictions = []\n\n    # 遍历每一行数据进行预测\n    for _, row in test_df.iterrows():\n        symbol_id = row['symbol_id']\n        row_id = row['row_id']\n        \n        # 加载对应 symbol_id 的模型\n        model = load_model(symbol_id)\n        \n        if model is not None:\n            try:\n                # 提取特征列并创建测试数据\n                X_test = pd.DataFrame([row])[feature_columns]\n                \n                # 进行预测\n                pred = model.predict(X_test)[0]\n                \n            except Exception as e:\n                print(f\"Error predicting for symbol {symbol_id}: {str(e)}\")\n                pred = 0.0\n        else:\n            pred = 0.0\n\n        # 记录预测结果\n        predictions.append({\n            'row_id': row_id,\n            'responder_6': pred\n        })\n    \n    # 将预测结果转换为 Polars DataFrame\n    predictions_df = pl.DataFrame(predictions)\n    \n    # 验证预测结果格式\n    assert list(predictions_df.columns) == ['row_id', 'responder_6']\n    assert len(predictions_df) == len(test)\n    \n    return predictions_df","metadata":{"execution":{"iopub.status.busy":"2024-10-27T08:53:56.510894Z","iopub.execute_input":"2024-10-27T08:53:56.511550Z","iopub.status.idle":"2024-10-27T08:53:57.184806Z","shell.execute_reply.started":"2024-10-27T08:53:56.511505Z","shell.execute_reply":"2024-10-27T08:53:57.183827Z"},"trusted":true},"execution_count":null,"outputs":[]},{"cell_type":"code","source":"def reduce_mem_usage(df, float16_as32=True):\n    \"\"\"Optimize memory usage of a DataFrame.\"\"\"\n    start_mem = df.memory_usage().sum() / 1024**2\n    print('Memory usage of dataframe is {:.2f} MB'.format(start_mem))\n\n    for col in df.columns:\n        col_type = df[col].dtype\n        if col_type != object and str(col_type) != 'category':\n            c_min, c_max = df[col].min(), df[col].max()\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 c_min > np.iinfo(np.int64).min and c_max < np.iinfo(np.int64).max:\n                    df[col] = df[col].astype(np.int64)\n            else:\n                if c_min > np.finfo(np.float16).min and c_max < np.finfo(np.float16).max:\n                    if float16_as32:\n                        df[col] = df[col].astype(np.float32)\n                    else:\n                        df[col] = df[col].astype(np.float16)\n                elif 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('Memory usage after optimization is: {:.2f} MB'.format(end_mem))\n    print('Decreased by {:.1f}%'.format(100 * (start_mem - end_mem) / start_mem))\n\n    return df","metadata":{},"execution_count":null,"outputs":[]},{"cell_type":"code","source":"import os\nimport lightgbm as lgb\nimport pandas as pd\nimport polars as pl\nimport numpy as np\nfrom sklearn.metrics import mean_squared_error\nfrom typing import Dict, Optional, List\n\n# Global variables\nmodels: Dict[int, lgb.Booster] = {}\nfeature_columns: List[str] = []  # 存储特征列名\n\ndef load_partition(partition_file: str) -> pl.DataFrame:\n    \"\"\"Load a single partition from the specified file.\"\"\"\n    print(f\"Loading partition: {partition_file}\")\n    partition_df = pl.read_parquet(partition_file).to_pandas()  # 转换为 pandas DataFrame 以便优化内存\n    partition_df = reduce_mem_usage(partition_df, float16_as32=False)\n    return pl.read_parquet(partition_file)\n\ndef prepare_features(df: pd.DataFrame, is_training: bool = True) -> tuple:\n    \"\"\"\n    Prepare features for training/prediction.\n    \n    Args:\n        df: Input DataFrame\n        is_training: If True, update feature_columns; if False, use existing feature_columns\n    Returns:\n        Tuple of feature matrix and target variable\n    \"\"\"\n    global feature_columns\n    \n    # Define columns to exclude\n    exclude_cols = ['symbol_id', 'date_id', 'responder_6', 'row_id', 'weight', 'time_id', 'is_scored']\n    \n    # 只使用 feature_XX 列作为特征\n    available_features = [col for col in df.columns if col.startswith('feature_')]\n    \n    if is_training:\n        # During training, update feature_columns\n        feature_columns = available_features\n        print(f\"Number of features selected: {len(feature_columns)}\")\n        print(\"Feature columns:\", feature_columns)\n    else:\n        # 确保测试数据使用相同的特征列\n        if not all(col in df.columns for col in feature_columns):\n            missing_cols = [col for col in feature_columns if col not in df.columns]\n            print(f\"Missing features in test data: {missing_cols}\")\n            raise ValueError(\"Test data missing some training features\")\n    \n    # Select features\n    X = df[feature_columns]\n    y = df['responder_6'] if 'responder_6' in df.columns else None\n    \n    return X, y\n\ndef evaluate_model(test: pl.DataFrame, lags: Optional[pl.DataFrame] = None) -> pd.DataFrame:\n    \"\"\"Evaluate models on test data and generate predictions.\"\"\"\n    test_df = test.to_pandas()\n    test_df = reduce_mem_usage(test_df, float16_as32=True) \n    predictions = []\n    \n    print(\"\\nPrediction phase:\")\n    print(f\"Number of features in stored feature list: {len(feature_columns)}\")\n    \n    # Process each row in test data\n    for _, row in test_df.iterrows():\n        symbol_id = row['symbol_id']\n        row_id = row['row_id']\n        \n        try:\n            if symbol_id in models:\n                # Prepare features using stored feature_columns (is_training=False)\n                X_test, _ = prepare_features(pd.DataFrame([row]), is_training=False)\n                \n                # Make prediction\n                pred = models[symbol_id].predict(X_test)[0]\n            else:\n                print(f\"No model found for symbol {symbol_id}, using default prediction\")\n                pred = 0.0\n                \n            predictions.append({\n                'row_id': row_id,\n                'responder_6': pred\n            })\n                \n        except Exception as e:\n            print(f\"Error predicting for row {row_id}: {str(e)}\")\n            print(f\"Symbol ID: {symbol_id}\")\n            predictions.append({\n                'row_id': row_id,\n                'responder_6': 0.0\n            })\n    \n    predictions_df = pd.DataFrame(predictions)\n    print(f\"\\nPredictions summary:\")\n    print(f\"Total predictions made: {len(predictions_df)}\")\n    print(f\"Non-zero predictions: {(predictions_df['responder_6'] != 0).sum()}\")\n    print(f\"Mean prediction value: {predictions_df['responder_6'].mean():.4f}\")\n    \n    return predictions_df\n\ndef train(train_data: pl.DataFrame) -> None:\n    \"\"\"Train LightGBM models for each financial instrument (symbol_id).\"\"\"\n    df = train_data.to_pandas()\n    df = reduce_mem_usage(df, float16_as32=True)\n    \n    # Print data info\n    print(\"\\nDataset Info:\")\n    print(f\"Total columns: {len(df.columns)}\")\n    print(f\"Total rows: {len(df)}\")\n    \n    params = {\n        'objective': 'regression',\n        'metric': 'rmse',\n        'boosting_type': 'gbdt',\n        'device_type': 'gpu',\n        'learning_rate': 0.05,\n        'num_leaves': 31,\n        'max_depth': 6,\n        'feature_fraction': 0.8,\n        'bagging_fraction': 0.8,\n        'bagging_freq': 5,\n        'min_data_in_leaf': 50,\n        'verbose': -1,\n        'seed': 42\n    }\n    \n    # Train a model for each symbol_id\n    for symbol in df['symbol_id'].unique():\n        symbol_data = df[df['symbol_id'] == symbol].copy()\n        \n        if len(symbol_data) < 100:\n            print(f\"Skipping symbol {symbol} due to insufficient data\")\n            continue\n            \n        try:\n            # Prepare features\n            X, y = prepare_features(symbol_data, is_training=True)\n            \n            # Create dataset\n            train_dataset = lgb.Dataset(X, label=y)\n            \n            # Train model\n            model = lgb.train(\n                params,\n                train_dataset,\n                num_boost_round=1000\n            )\n            \n            # Store model\n            models[symbol] = model\n            print(f\"Successfully trained LightGBM model for symbol {symbol}\")\n            \n            model_path = f\"model_{symbol}.joblib\"\n            joblib.dump(models[symbol], model_path)\n            print(f\"Model for symbol {symbol} saved at {model_path}\")\n            \n        except Exception as e:\n            print(f\"Error training model for symbol {symbol}: {str(e)}\")\n    \n    print(f\"\\nTraining complete. Models trained for {len(models)} symbols.\")\n    print(f\"Models available for symbols: {list(models.keys())}\")\n\ndef weighted_r2_score(y_true: np.ndarray, y_pred: np.ndarray, weights: np.ndarray) -> float:\n    \"\"\"Calculate weighted R² score.\"\"\"\n    weighted_mean = np.sum(weights * y_true) / np.sum(weights)\n    weighted_ss_tot = np.sum(weights * (y_true - weighted_mean) ** 2)\n    weighted_ss_res = np.sum(weights * (y_true - y_pred) ** 2)\n    \n    return 1 - (weighted_ss_res / weighted_ss_tot)","metadata":{"execution":{"iopub.status.busy":"2024-10-27T08:53:57.186604Z","iopub.execute_input":"2024-10-27T08:53:57.187127Z","iopub.status.idle":"2024-10-27T08:54:00.749561Z","shell.execute_reply.started":"2024-10-27T08:53:57.187086Z","shell.execute_reply":"2024-10-27T08:54:00.748627Z"},"trusted":true},"execution_count":null,"outputs":[]},{"cell_type":"code","source":"partition_folder = '/kaggle/input/jane-street-real-time-market-data-forecasting/train.parquet'\npartition_files = [os.path.join(partition_folder, f\"partition_id={i}\") for i in range(1)]\ntest_data = pl.read_parquet('/kaggle/input/jane-street-real-time-market-data-forecasting/test.parquet')\ntest_data = test_data.to_pandas()\ntest_data = reduce_mem_usage(test_data, float16_as32=False)  # 优化内存\ntest_data = pl.from_pandas(test_data) \nlags_data = None\n\nis_train = True\n\nif is_train:\n\n    # Train models\n    for partition_file in partition_files:\n        if os.path.exists(partition_file):\n            partition_data = load_partition(partition_file)\n            train(partition_data)\n        else:\n            print(f\"Partition not found: {partition_file}\")\n\n    # Generate predictions\n    predictions = evaluate_model(test_data, lags_data)","metadata":{"execution":{"iopub.status.busy":"2024-10-27T08:54:00.751542Z","iopub.execute_input":"2024-10-27T08:54:00.752849Z","iopub.status.idle":"2024-10-27T08:58:13.689870Z","shell.execute_reply.started":"2024-10-27T08:54:00.752789Z","shell.execute_reply":"2024-10-27T08:58:13.688720Z"},"trusted":true},"execution_count":null,"outputs":[]},{"cell_type":"code","source":"import kaggle_evaluation\ninference_server = kaggle_evaluation.jane_street_inference_server.JSInferenceServer(predict)\n\nif os.getenv('KAGGLE_IS_COMPETITION_RERUN'):\n    inference_server.serve()\nelse:\n    inference_server.run_local_gateway(\n        (\n            '/kaggle/input/jane-street-real-time-market-data-forecasting/test.parquet',\n            '/kaggle/input/jane-street-real-time-market-data-forecasting/lags.parquet',\n        )\n    )","metadata":{"execution":{"iopub.status.busy":"2024-10-27T09:00:34.618528Z","iopub.execute_input":"2024-10-27T09:00:34.618977Z","iopub.status.idle":"2024-10-27T09:00:35.574285Z","shell.execute_reply.started":"2024-10-27T09:00:34.618938Z","shell.execute_reply":"2024-10-27T09:00:35.569781Z"},"trusted":true},"execution_count":null,"outputs":[]},{"cell_type":"code","source":"","metadata":{},"execution_count":null,"outputs":[]}]}