{"metadata":{"kernelspec":{"name":"python3","display_name":"Python 3","language":"python"},"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":"nvidiaTeslaT4","dataSources":[{"sourceId":84493,"databundleVersionId":9871156,"sourceType":"competition"}],"dockerImageVersionId":30787,"isInternetEnabled":true,"language":"python","sourceType":"notebook","isGpuEnabled":true}},"nbformat_minor":4,"nbformat":4,"cells":[{"cell_type":"code","source":"import os\nimport joblib \nimport gc\nimport pandas as pd\nimport polars as pl\n#import lightgbm as lgb \n# import xgboost as xgb\n# import catboost as cbt\nimport numpy as np \nimport torch\n\nfrom joblib import Parallel, delayed\n\nimport kaggle_evaluation.jane_street_inference_server\n\n# apt-get update\n# apt-get install --no-install-recommends libnvidia-compute-535 ocl-icd-opencl-dev\n# apt-get install --no-install-recommends opencl-headers ocl-icd-libopencl1 ocl-icd-opencl-dev\n# !pip install grpcio pyarrow protobuf pandas scikit-learn joblib\n# !pip install lightgbm==4.2.0 -i https://mirrors.aliyun.com/pypi/simple/\n# !pip install catboost==1.2.7 -i https://mirrors.aliyun.com/pypi/simple/\n# !pip install xgboost==2.0.3 -i https://mirrors.aliyun.com/pypi/simple/\n# !pip install joblib==1.4.2 -i https://mirrors.aliyun.com/pypi/simple/\n\n\ndef reduce_mem_usage(self, float16_as32=True):\n    #memory_usage()是df每列的内存使用量,sum是对它们求和, B->KB->MB\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#列名的type\n        if col_type != object and str(col_type)!='category':#不是object也就是说这里处理的是数值类型的变量\n            c_min,c_max = df[col].min(),df[col].max() #求出这列的最大值和最小值\n            if str(col_type)[:3] == 'int':#如果是int类型的变量,不管是int8,int16,int32还是int64\n                #如果这列的取值范围是在int8的取值范围内,那就对类型进行转换 (-128 到 127)\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                #如果这列的取值范围是在int16的取值范围内,那就对类型进行转换(-32,768 到 32,767)\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                #如果这列的取值范围是在int32的取值范围内,那就对类型进行转换(-2,147,483,648到2,147,483,647)\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                #如果这列的取值范围是在int64的取值范围内,那就对类型进行转换(-9,223,372,036,854,775,808到9,223,372,036,854,775,807)\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                #如果数值在float16的取值范围内,如果觉得需要更高精度可以考虑float32\n                if c_min > np.finfo(np.float16).min and c_max < np.finfo(np.float16).max:\n                    if float16_as32:#如果数据需要更高的精度可以选择float32\n                        df[col] = df[col].astype(np.float32)\n                    else:\n                        df[col] = df[col].astype(np.float16)  \n                #如果数值在float32的取值范围内，对它进行类型转换\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                #如果数值在float64的取值范围内，对它进行类型转换\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    #相比一开始的内存减少了百分之多少\n    print('Decreased by {:.1f}%'.format(100 * (start_mem - end_mem) / start_mem))\n\n    return df","metadata":{"execution":{"iopub.status.busy":"2024-10-24T09:15:40.926424Z","iopub.execute_input":"2024-10-24T09:15:40.926796Z","iopub.status.idle":"2024-10-24T09:15:45.476857Z","shell.execute_reply.started":"2024-10-24T09:15:40.926761Z","shell.execute_reply":"2024-10-24T09:15:45.476055Z"},"trusted":true},"execution_count":null,"outputs":[]},{"cell_type":"code","source":"# Define the path to the input data directory\n# If the local directory exists, use it; otherwise, use the Kaggle input directory\ninput_path = '/kaggle/input/jane-street-real-time-market-data-forecasting/' #if os.path.exists('/workspace/jane/jane-street-real-time-market-data-forecasting') else '/kaggle/input/jane-street-real-time-market-data-forecasting/'\n\n# Flag to determine if the script is in training mode or not\nTRAINING = True\n\n# Define the feature names based on the number of features (79 in this case)\nfeature_names = [f\"feature_{i:02d}\" for i in range(79)] \n\n# Number of validation dates to use\nnum_valid_dates = 100\n\n# Number of dates to skip from the beginning of the dataset\nskip_dates = 500\n\n# Number of folds for cross-validation\nN_fold = 5","metadata":{"execution":{"iopub.status.busy":"2024-10-24T09:15:45.478651Z","iopub.execute_input":"2024-10-24T09:15:45.479151Z","iopub.status.idle":"2024-10-24T09:15:45.484773Z","shell.execute_reply.started":"2024-10-24T09:15:45.479108Z","shell.execute_reply":"2024-10-24T09:15:45.483928Z"},"trusted":true},"execution_count":null,"outputs":[]},{"cell_type":"code","source":"# If in training mode, load the training data\nif TRAINING:\n    # Load the training data from a Parquet file\n    df = pd.read_parquet(f'{input_path}/train.parquet')\n    \n    # Reduce memory usage of the DataFrame (function not provided here)\n    df = reduce_mem_usage(df, False)\n    \n    # Filter the DataFrame to include only dates greater than or equal to skip_dates\n    df = df[df['date_id'] >= skip_dates].reset_index(drop=True)\n    \n    # Get unique dates from the DataFrame\n    dates = df['date_id'].unique()\n    \n    # Define validation dates as the last `num_valid_dates` dates\n    valid_dates = dates[-num_valid_dates:]\n    \n    # Define training dates as all dates except the last `num_valid_dates` dates\n    train_dates = dates[:-num_valid_dates]\n    \n    # Display the last few rows of the DataFrame (for debugging purposes)\n    print(df.tail())","metadata":{"execution":{"iopub.status.busy":"2024-10-24T09:15:45.486091Z","iopub.execute_input":"2024-10-24T09:15:45.486534Z"},"trusted":true},"execution_count":null,"outputs":[]},{"cell_type":"markdown","source":"# Split","metadata":{}},{"cell_type":"code","source":"selected_dates = [date for ii, date in enumerate(train_dates)]\nX_train = df[feature_names].loc[df['date_id'].isin(selected_dates)]\ny_train = df['responder_6'].loc[df['date_id'].isin(selected_dates)]\nw_train = df['weight'].loc[df['date_id'].isin(selected_dates)]","metadata":{"trusted":true},"execution_count":null,"outputs":[]},{"cell_type":"code","source":"# Extract features, target, and weights for validation dates\nvalid_dates = valid_dates.tolist()\nX_valid = df[feature_names].loc[df['date_id'].isin(valid_dates)]\ny_valid = df['responder_6'].loc[df['date_id'].isin(valid_dates)]\nw_valid = df['weight'].loc[df['date_id'].isin(valid_dates)]","metadata":{"trusted":true},"execution_count":null,"outputs":[]},{"cell_type":"code","source":"del df\ngc.collect()","metadata":{"trusted":true},"execution_count":null,"outputs":[]},{"cell_type":"markdown","source":"# Missing indicator","metadata":{}},{"cell_type":"code","source":"from sklearn.impute import MissingIndicator\n\n# Assume your original DataFrame is named 'df'\n# Create a MissingIndicator instance\nindicator = MissingIndicator(features='all', error_on_new=False)\n\n# Fit the indicator to the DataFrame\nindicator_matrix_train = indicator.fit_transform(X_train)\nindicator_matrix_valid = indicator.transform(X_valid)","metadata":{"trusted":true},"execution_count":null,"outputs":[]},{"cell_type":"code","source":"indicator_matrix_train.shape, indicator_matrix_valid.shape","metadata":{"trusted":true},"execution_count":null,"outputs":[]},{"cell_type":"markdown","source":"# Imputer","metadata":{}},{"cell_type":"code","source":"from sklearn.impute import SimpleImputer\nimp_mean = SimpleImputer(missing_values=np.nan, strategy='mean')\nimputed_X_train = imp_mean.fit_transform(X_train)\nimputed_X_valid = imp_mean.transform(X_valid)","metadata":{"trusted":true},"execution_count":null,"outputs":[]},{"cell_type":"code","source":"imputed_X_train.shape, imputed_X_valid.shape","metadata":{"trusted":true},"execution_count":null,"outputs":[]},{"cell_type":"markdown","source":"# Scaling","metadata":{}},{"cell_type":"code","source":"from sklearn.preprocessing import RobustScaler\n# Initialize and fit the RobustScaler on the training features\nscaler = RobustScaler()\n\n# If not in training mode, load the pre-trained model and scaler\n# scaler = joblib.load(f'scalers/robust_scaler.pkl')\nX_train_scaled = scaler.fit_transform(imputed_X_train)\nX_valid_scaled = scaler.transform(imputed_X_valid)","metadata":{"trusted":true},"execution_count":null,"outputs":[]},{"cell_type":"code","source":"X_train_scaled.shape, X_valid_scaled.shape","metadata":{"trusted":true},"execution_count":null,"outputs":[]},{"cell_type":"code","source":"joblib.dump(indicator, f'indicators/indicator.pkl')\njoblib.dump(imp_mean, f'imputers/imputer.pkl')\njoblib.dump(scaler, f'scalers/robust_scaler.pkl')","metadata":{"trusted":true},"execution_count":null,"outputs":[]},{"cell_type":"code","source":"indicator =  joblib.load(f'indicators/indicator.pkl')\nimp_mean = joblib.load(f'imputers/imputer.pkl')\nscaler = joblib.load(f'scalers/robust_scaler.pkl')","metadata":{"trusted":true},"execution_count":null,"outputs":[]},{"cell_type":"markdown","source":"# Final dataset","metadata":{}},{"cell_type":"code","source":"del X_train, X_valid\ngc.collect()","metadata":{"trusted":true},"execution_count":null,"outputs":[]},{"cell_type":"code","source":"X_train = np.concatenate([X_train_scaled, indicator_matrix_train[:, 4:]], 1)\nX_valid = np.concatenate([X_valid_scaled, indicator_matrix_valid[:, 4:]], 1)","metadata":{"trusted":true},"execution_count":null,"outputs":[]},{"cell_type":"code","source":"# X_train.shape, X_valid.shape","metadata":{"trusted":true},"execution_count":null,"outputs":[]},{"cell_type":"code","source":"joblib.dump(X_train, f'datas/X_train.pkl')\njoblib.dump(X_valid, f'datas/X_valid.pkl')\njoblib.dump(y_train, f'datas/y_train.pkl')\njoblib.dump(y_valid, f'datas/y_valid.pkl')\njoblib.dump(w_train, f'datas/w_train.pkl')\njoblib.dump(w_valid, f'datas/w_valid.pkl')","metadata":{"trusted":true},"execution_count":null,"outputs":[]},{"cell_type":"code","source":"del X_train_scaled, X_valid_scaled, indicator_matrix_train, indicator_matrix_valid, imputed_X_train, imputed_X_valid\ngc.collect()","metadata":{},"execution_count":null,"outputs":[]},{"cell_type":"markdown","source":"# Load data","metadata":{}},{"cell_type":"code","source":"X_train = joblib.load(f'datas/X_train.pkl')\nX_valid = joblib.load(f'datas/X_valid.pkl')\ny_train = joblib.load(f'datas/y_train.pkl')\ny_valid = joblib.load(f'datas/y_valid.pkl')\nw_train = joblib.load(f'datas/w_train.pkl')\nw_valid = joblib.load(f'datas/w_valid.pkl')","metadata":{},"execution_count":null,"outputs":[]},{"cell_type":"markdown","source":"# Model Training","metadata":{}},{"cell_type":"code","source":"import os\nfrom torch import optim, nn, utils, Tensor\nfrom torch.utils.data import Dataset, DataLoader\nfrom torchvision.transforms import ToTensor\nimport lightning as L\n\nclass MLP_model(nn.Module):\n    def __init__(self, d_model=256) -> None:\n        super().__init__()      \n        self.d_model = d_model\n        self.linear1 = nn.Linear(154, d_model)\n        self.linear2 = nn.Linear(d_model, 1)\n        self.slinear1 = nn.Linear(d_model, d_model)\n        self.slinear2 = nn.Linear(d_model, d_model)\n        self.act = nn.SiLU()\n        \n           \n    def forward(self, src) -> Tensor:\n        src = self.linear1(src)\n        u = self.slinear1(src)\n        v = self.act(self.slinear2(src))\n        output = u * v\n        output = self.linear2(output)\n        return output\n\nmodel = MLP_model()","metadata":{},"execution_count":null,"outputs":[]},{"cell_type":"code","source":"from weighted_r2 import R2Score","metadata":{},"execution_count":null,"outputs":[]},{"cell_type":"code","source":"# Define a dataset\n\nclass CustomDataset(Dataset):\n    def __init__(self, features, targets, weights):\n        self.features = features\n        self.targets = targets\n        self.weights = weights\n\n    def __len__(self):\n        return len(self.features)\n\n    def __getitem__(self, idx):\n        x = self.features[idx]\n        y = self.targets.iloc[idx]\n        w = self.weights.iloc[idx] \n        return torch.tensor(x, dtype=torch.float32), torch.tensor(y, dtype=torch.float32), torch.tensor(w, dtype=torch.float32)\n    \n# Assume X_train_scaled and y_train are numpy arrays\ntrain_dataset = CustomDataset(X_train, y_train, w_train)\nvalid_dataset = CustomDataset(X_valid, y_valid, w_valid)\n\n\n# # Create a DataLoader with desired batch size\nbatch_size = 256\ntrain_loader = DataLoader(train_dataset, batch_size=batch_size, shuffle=True)\nvalid_loader = DataLoader(valid_dataset, batch_size=batch_size*4)","metadata":{},"execution_count":null,"outputs":[]},{"cell_type":"code","source":"# define the LightningModule\nclass LitAutoEncoder(L.LightningModule):\n    def __init__(self, encoder):\n        super().__init__()\n        self.model = model\n        self.r2score = R2Score()\n        #self.lr = learning_rate\n\n    def training_step(self, batch, batch_idx):\n        # training_step defines the train loop.\n        # it is independent of forward\n        x, y, w = batch\n        y_hat = self.model(x).flatten()\n        loss = nn.functional.mse_loss(y_hat, y)\n        # Logging to TensorBoard (if installed) by default\n        self.log(\"train_loss\", loss)\n        return loss\n    \n    def validation_step(self, batch, batch_idx):\n        x, y, w = batch\n        y_hat = self.model(x).flatten()\n        loss = nn.functional.mse_loss(y_hat, y)\n        self.log(\"val_loss\", loss, on_step=False, on_epoch=True)\n        self.r2score(y_hat, y, w)\n        self.log(\"val_r2\", self.r2score, on_step=False, on_epoch=True)\n\n    def training_step(self, batch, batch_idx):\n        # training_step defines the train loop.\n        # it is independent of forward\n        x, y, w = batch\n        y_hat = self.model(x).flatten()\n        loss = nn.functional.mse_loss(y_hat, y)\n        # Logging to TensorBoard (if installed) by default\n        self.log(\"train_loss\", loss)\n        return loss\n\n    def configure_optimizers(self):\n        optimizer = optim.AdamW(self.parameters(), lr=1e-3)\n        return optimizer\n    \n    def forward(self, inputs):\n        return self.model(inputs)\n    \n    def predict_step(self, batch, batch_idx, dataloader_idx=0):\n        y_hat = self.forward(batch).flatten()\n        return y_hat\n    \n# init the network\nnet = LitAutoEncoder(model)","metadata":{},"execution_count":null,"outputs":[]},{"cell_type":"markdown","source":"# Learning rate finder","metadata":{}},{"cell_type":"code","source":"# from lightning.pytorch.tuner import Tuner\n# # Create a Tuner\n# trainer = L.Trainer(max_epochs=1,)\n# tuner = Tuner(trainer)\n\n# # finds learning rate automatically\n# # sets hparams.lr or hparams.learning_rate to that learning rate\n# lr_finder = tuner.lr_find(net, attr_name='lr')\n# # Plot with\n# fig = lr_finder.plot(suggest=True)\n# fig.show()","metadata":{},"execution_count":null,"outputs":[]},{"cell_type":"markdown","source":"# Training","metadata":{}},{"cell_type":"code","source":"from lightning.pytorch.callbacks.early_stopping import EarlyStopping\nfrom lightning.pytorch.callbacks import ModelCheckpoint\n\nearlystopping = EarlyStopping(monitor=\"val_r2\", mode=\"max\", patience=7)\n\ncheckpoint_callback = ModelCheckpoint(\n    monitor='val_r2',   # Metric to monitor\n    dirpath='checkpoints/',  # Directory to save the checkpoints\n    filename='MLP-checkpoint-{val_r2:.2f}',  # Checkpoint file name\n    save_top_k=1,          # Save the top 1 checkpoint\n    mode='max',            # Minimize the monitored metric (for val_loss)\n    verbose=True           # Print information during training\n)","metadata":{},"execution_count":null,"outputs":[]},{"cell_type":"code","source":"# train the model (hint: here are some helpful Trainer arguments for rapid idea iteration)\n#trainer = L.Trainer(limit_train_batches=100, max_epochs=1, limit_val_batches=100)\n#trainer = L.Trainer(limit_train_batches=100, max_epochs=10, limit_val_batches=100, callbacks=[earlystopping, checkpoint_callback])\ntrainer = L.Trainer(max_epochs=100, callbacks=[earlystopping, checkpoint_callback])\ntrainer.fit(model=net, train_dataloaders=train_loader, val_dataloaders=valid_loader)","metadata":{},"execution_count":null,"outputs":[]},{"cell_type":"markdown","source":"# Predicting","metadata":{}},{"cell_type":"code","source":"# for predicting data\nindicator =  joblib.load(f'indicators/indicator.pkl')\nimp_mean = joblib.load(f'imputers/imputer.pkl')\nscaler = joblib.load(f'scalers/robust_scaler.pkl')\n\nclass PredDataset(Dataset):\n    def __init__(self, features):\n        self.features = features\n\n    def __len__(self):\n        return len(self.features)\n\n    def __getitem__(self, idx):\n        x = self.features[idx]\n        return torch.tensor(x, dtype=torch.float32)\n    \n# # Create a DataLoader with desired batch size\nbatch_size = 128\ntrain_loader = DataLoader(train_dataset, batch_size=batch_size, shuffle=True)\n\ndef transform_pred(test):   \n    feat = test[feature_names].to_numpy()\n    indicator_matrix_test = indicator.transform(feat)\n    imputed_X_test = imp_mean.transform(feat)\n\n    # Check for NaN in the imputed data\n    if np.any(np.isnan(imputed_X_test)):\n        print(\"NaN detected in imputed_X_test\")\n    else:\n        print(\"No NaN in imputed_X_test\")\n\n    X_test_scaled = scaler.transform(imputed_X_test)\n    X_test = np.concatenate([X_test_scaled, indicator_matrix_test[:, 4:]], 1)\n    pred_dataset = PredDataset(X_test)\n    pred_loader = DataLoader(pred_dataset, batch_size=128)\n    return pred_loader","metadata":{},"execution_count":null,"outputs":[]},{"cell_type":"code","source":"lags_ : pl.DataFrame | None = None\n\n# Replace this function with your inference code.\n# You can return either a Pandas or Polars dataframe, though Polars is recommended.\n# Each batch of predictions (except the very first) must be returned within 10 minutes of the batch features being provided.\ndef predict(test: pl.DataFrame, lags: pl.DataFrame | None) -> pl.DataFrame | pd.DataFrame:\n    \"\"\"Make a prediction.\"\"\"\n    # All the responders from the previous day are passed in at time_id == 0. We save them in a global variable for access at every time_id.\n    # Use them as extra features, if you like.\n    global lags_\n    if lags is not None:\n        lags_ = lags\n\n    predictions = test.select(\n        'row_id',\n        pl.lit(0.0).alias('responder_6'),\n    )\n    \n    pred_loader = transform_pred(test)\n    pred = trainer.predict(net, pred_loader)\n    pred = torch.cat(pred, 0).numpy()\n    predictions = predictions.with_columns(pl.Series('responder_6', pred))\n\n    # The predict function must return a DataFrame\n    assert isinstance(predictions, pl.DataFrame | pd.DataFrame)\n    # with columns 'row_id', 'responer_6'\n    assert list(predictions.columns) == ['row_id', 'responder_6']\n    # and as many rows as the test data.\n    assert len(predictions) == len(test)\n    #print(predictions)\n    return predictions","metadata":{},"execution_count":null,"outputs":[]},{"cell_type":"code","source":"inference_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            '/workspace/jane/jane-street-real-time-market-data-forecasting/test.parquet',\n            '/workspace/jane/jane-street-real-time-market-data-forecasting/lags.parquet',\n        )\n    )","metadata":{},"execution_count":null,"outputs":[]},{"cell_type":"code","source":"","metadata":{},"execution_count":null,"outputs":[]},{"cell_type":"code","source":"","metadata":{},"execution_count":null,"outputs":[]}]}