{"metadata":{"kernelspec":{"language":"python","display_name":"Python 3","name":"python3"},"language_info":{"pygments_lexer":"ipython3","nbconvert_exporter":"python","version":"3.6.4","file_extension":".py","codemirror_mode":{"name":"ipython","version":3},"name":"python","mimetype":"text/x-python"}},"nbformat_minor":4,"nbformat":4,"cells":[{"cell_type":"markdown","source":"This is an implementation of the [HIST model](https://arxiv.org/abs/2110.13716) model, it was designed for stock trend forecasting. The concept was the sector (i.e. technology, internet retail, etc) for a particular stock. For the AMEX challenge, I use K mean clustering to separate the the P R S B columns in to different cluster and one hot encoded it as a concept. \n\nThe model is also a time series model and therefore I used the [RNN / transformer dataset](https://www.kaggle.com/code/cdeotte/tensorflow-gru-starter-0-790) from Chris Deotte.\n\nThe performance is not good, not even close to a Transformer / GRU.\n\nI also use PyTorch Lightning for the training. ","metadata":{}},{"cell_type":"code","source":"import sys\nsys.path.insert(1, '../input/stock2concept/')\nimport warnings\nwarnings.filterwarnings(\"ignore\")\n\nimport math\nimport torch\nfrom torch.utils.data import DataLoader\nimport torch.nn.functional as F\nimport torch.nn as nn\nfrom torch.optim.lr_scheduler import CosineAnnealingWarmRestarts, ExponentialLR, CosineAnnealingLR, StepLR, OneCycleLR\n\nfrom model_hist import HISTModel\nimport pytorch_lightning as pl\nfrom pytorch_lightning.loggers import WandbLogger\nfrom pytorch_lightning.callbacks import ModelCheckpoint\n\nimport pandas as pd","metadata":{"_uuid":"8f2839f25d086af736a60e9eeb907d3b93b6e0e5","_cell_guid":"b1076dfc-b9ad-4769-8c92-a6c4dae69d19","execution":{"iopub.status.busy":"2022-07-15T07:36:59.900748Z","iopub.execute_input":"2022-07-15T07:36:59.901193Z","iopub.status.idle":"2022-07-15T07:37:04.771897Z","shell.execute_reply.started":"2022-07-15T07:36:59.901155Z","shell.execute_reply":"2022-07-15T07:37:04.770860Z"},"trusted":true},"execution_count":null,"outputs":[]},{"cell_type":"markdown","source":"# Dataset","metadata":{}},{"cell_type":"code","source":"import pandas as pd\nfrom torch.utils.data import Dataset, DataLoader\nimport numpy as np\nfrom tqdm import tqdm \n\nclass Dataset_AMEX(Dataset):\n    def __init__(self, flag='train', fold=1):\n        assert flag in ['train', 'val', 'test']\n        if flag in ['train', 'val']:\n            assert fold in range(10)       \n        self.PATH_TO_DATA = '../input/amex-data-for-transformers-and-rnns/data/'\n        self.flag = flag\n        self.fold = fold\n        self.__read_data__()\n        \n        \n    def __read_data__(self):\n        valid_idx = [2*self.fold+1, 2*self.fold+2]\n        train_idx = [x for x in [1,2,3,4,5,6,7,8,9,10] if x not in valid_idx]\n        test_idx = [1,2,3,4,5,6,7,8,9,10,11,12,13,14,15,16,17,18,19,20]\n        if self.flag == 'train':\n            X_train = []; y_train = []; s2c_train = []\n            for k in train_idx:\n                X_train.append( np.load(f'{self.PATH_TO_DATA}data_{k}.npy'))\n                y_train.append( pd.read_parquet(f'{self.PATH_TO_DATA}targets_{k}.pqt') )\n                s2c_train.append(np.load(f'../input/stock2concept/stock2concept/s2c_{k}.npy', allow_pickle=True))\n            self.X = np.concatenate(X_train,axis=0)\n            self.y = pd.concat(y_train).target.values\n            self.s2c = np.concatenate(s2c_train,axis=0).astype(np.int)\n            del X_train, y_train, s2c_train\n            print('### Training data shapes', self.X.shape, self.y.shape)\n        elif self.flag == 'val':\n            X_valid = []; y_valid = []; s2c_valid = []\n            for k in valid_idx:\n                X_valid.append(np.load(f'{self.PATH_TO_DATA}data_{k}.npy'))\n                y_valid.append( pd.read_parquet(f'{self.PATH_TO_DATA}targets_{k}.pqt') )\n                s2c_valid.append(np.load(f'../input/stock2concept/stock2concept/s2c_{k}.npy', allow_pickle=True))\n            self.X = np.concatenate(X_valid,axis=0)\n            self.y = pd.concat(y_valid).target.values\n            self.s2c = np.concatenate(s2c_valid,axis=0).astype(np.int)\n            del X_valid, y_valid, s2c_valid\n            print('### Validation data shapes', self.X.shape, self.y.shape)\n        elif self.flag == 'test':\n            X_test = []; y_test = []; s2c_test = [] \n            for k in tqdm(test_idx):\n                X_test.append(np.load(f'{self.PATH_TO_DATA}test_data_{k}.npy'))\n                y_test.append(np.zeros(len(X_test)))\n                s2c_test.append(np.load(f'../input/stock2concept/stock2concept/s2c_test_{k}.npy',  allow_pickle=True))\n            self.X = np.concatenate(X_test,axis=0)\n            self.y = np.concatenate(y_test,axis=0)\n            self.s2c = np.concatenate(s2c_test,axis=0).astype(np.int)\n            del X_test, y_test, s2c_test\n            print('### Test data shapes', self.X.shape, self.y.shape)\n                              \n            #k=self.fold\n            #self.X = np.load(f'{self.PATH_TO_DATA}test_data_{k}.npy')\n            #self.y = np.zeros(len(self.X))\n            #print('### Test data shapes', self.X.shape, self.y.shape)\n    \n    def __len__(self):\n        return len(self.X)\n\n    def __getitem__(self, index):\n        if self.flag == 'test':\n            return self.X[index], self.s2c[index], np.empty_like(self.X[index])\n        else:\n            return self.X[index], self.s2c[index], self.y[index]","metadata":{"jupyter":{"source_hidden":true},"execution":{"iopub.status.busy":"2022-07-15T07:37:04.775758Z","iopub.execute_input":"2022-07-15T07:37:04.777315Z","iopub.status.idle":"2022-07-15T07:37:04.798496Z","shell.execute_reply.started":"2022-07-15T07:37:04.777277Z","shell.execute_reply":"2022-07-15T07:37:04.797574Z"},"trusted":true},"execution_count":null,"outputs":[]},{"cell_type":"markdown","source":"# Metric","metadata":{}},{"cell_type":"code","source":"import torch.nn.functional as F\nfrom torch import nn, Tensor\nfrom torchmetrics.utilities import rank_zero_warn\nimport pandas as pd\nimport gc\nimport numpy as np\nfrom tqdm.notebook import tqdm\n# Typing \nfrom typing import Optional\nfrom torchmetrics import Metric\nimport torch\n\nclass AmexMetric(Metric):\n    is_differentiable: Optional[bool] = False\n\n    # Set to True if the metric reaches it optimal value when the metric is maximized.\n    # Set to False if it when the metric is minimized.\n    higher_is_better: Optional[bool] = True\n\n    # Set to True if the metric during 'update' requires access to the global metric\n    # state for its calculations. If not, setting this to False indicates that all\n    # batch states are independent and we will optimize the runtime of 'forward'\n    full_state_update: bool = True\n\n    def __init__(self):\n        super().__init__()\n        \n        self.add_state(\"all_true\", default=[], dist_reduce_fx=\"cat\")\n        self.add_state(\"all_pred\", default=[], dist_reduce_fx=\"cat\")\n\n        rank_zero_warn(\n            \"Metric `Amex` will save all targets and predictions in buffer.\"\n            \" For large datasets this may lead to large memory footprint.\"\n        )\n\n    def update(self, y_pred: torch.Tensor, y_true: torch.Tensor):\n        \n        y_true = y_true.double()\n        y_pred = y_pred.double()\n        \n        self.all_true.append(y_true)\n        self.all_pred.append(y_pred)\n        \n    def compute(self):\n        y_true = torch.cat(self.all_true)\n        y_pred = torch.cat(self.all_pred)\n        # count of positives and negatives\n        n_pos = y_true.sum()\n        n_neg = y_pred.shape[0] - n_pos\n\n        # sorting by descring prediction values\n        indices = torch.argsort(y_pred, dim=0, descending=True)\n        preds, target = y_pred[indices], y_true[indices]\n\n        # filter the top 4% by cumulative row weights\n        weight = 20.0 - target * 19.0\n        cum_norm_weight = (weight / weight.sum()).cumsum(dim=0)\n        four_pct_filter = cum_norm_weight <= 0.04\n\n        # default rate captured at 4%\n        d = target[four_pct_filter].sum() / n_pos\n\n        # weighted gini coefficient\n        lorentz = (target / n_pos).cumsum(dim=0)\n        gini = ((lorentz - cum_norm_weight) * weight).sum()\n\n        # max weighted gini coefficient\n        gini_max = 10 * n_neg * (1 - 19 / (n_pos + 20 * n_neg))\n\n        # normalized weighted gini coefficient\n        g = gini / gini_max\n        \n        return 0.5 * (g + d)\n\n\ndef amex_metric(y_true: pd.DataFrame, y_pred: pd.DataFrame) -> float:\n    def top_four_percent_captured(y_true: pd.DataFrame, y_pred: pd.DataFrame) -> float:\n        df = (pd.concat([y_true, y_pred], axis='columns')\n              .sort_values('prediction', ascending=False))\n        df['weight'] = df['target'].apply(lambda x: 20 if x == 0 else 1)\n        four_pct_cutoff = int(0.04 * df['weight'].sum())\n        df['weight_cumsum'] = df['weight'].cumsum()\n        df_cutoff = df.loc[df['weight_cumsum'] <= four_pct_cutoff]\n        return (df_cutoff['target'] == 1).sum() / (df['target'] == 1).sum()\n\n    def weighted_gini(y_true: pd.DataFrame, y_pred: pd.DataFrame) -> float:\n        df = (pd.concat([y_true, y_pred], axis='columns')\n              .sort_values('prediction', ascending=False))\n        df['weight'] = df['target'].apply(lambda x: 20 if x == 0 else 1)\n        df['random'] = (df['weight'] / df['weight'].sum()).cumsum()\n        total_pos = (df['target'] * df['weight']).sum()\n        df['cum_pos_found'] = (df['target'] * df['weight']).cumsum()\n        df['lorentz'] = df['cum_pos_found'] / total_pos\n        df['gini'] = (df['lorentz'] - df['random']) * df['weight']\n        return df['gini'].sum()\n\n    def normalized_weighted_gini(y_true: pd.DataFrame, y_pred: pd.DataFrame) -> float:\n        y_true_pred = y_true.rename(columns={'target': 'prediction'})\n        return weighted_gini(y_true, y_pred) / weighted_gini(y_true, y_true_pred)\n\n    g = normalized_weighted_gini(y_true, y_pred)\n    d = top_four_percent_captured(y_true, y_pred)\n\n    return 0.5 * (g + d)","metadata":{"jupyter":{"source_hidden":true},"execution":{"iopub.status.busy":"2022-07-15T07:37:04.801517Z","iopub.execute_input":"2022-07-15T07:37:04.802155Z","iopub.status.idle":"2022-07-15T07:37:04.826415Z","shell.execute_reply.started":"2022-07-15T07:37:04.802117Z","shell.execute_reply":"2022-07-15T07:37:04.825469Z"},"trusted":true},"execution_count":null,"outputs":[]},{"cell_type":"code","source":"class Dataset_pl(pl.LightningDataModule):\n    def __init__(self, fold):\n        super().__init__()\n        self.fold = fold\n        \n    def prepare_data(self):\n        pass\n\n    def setup(self, stage= None):\n        # Assign train/val datasets for use in dataloaders\n        if stage == \"fit\" or stage is None:\n            self.train_set = Dataset_AMEX('train', fold=self.fold)\n            self.val_set = Dataset_AMEX('val', fold=self.fold)\n        if stage == \"validate\":\n            self.val_set = Dataset_AMEX('val', fold=self.fold)\n        # Assign test dataset for use in dataloader(s)\n        if stage == \"test\" or stage is None:\n            self.val_set = Dataset_AMEX('val', fold=self.fold)\n        if stage == \"predict\" or stage is None:\n            self.test_set = Dataset_AMEX('test')\n\n    def train_dataloader(self):\n        return DataLoader(self.train_set, batch_size=512, shuffle=True, num_workers=4)\n\n    def val_dataloader(self):\n        return DataLoader(self.val_set, batch_size=2048, shuffle=False, num_workers=4)\n\n    def test_dataloader(self):\n        return DataLoader(self.test_set, batch_size=4096, shuffle=False, num_workers=4)\n\n    def predict_dataloader(self):\n        return DataLoader(self.test_set, batch_size=4096, shuffle=False, num_workers=4)","metadata":{"execution":{"iopub.status.busy":"2022-07-15T07:37:04.828586Z","iopub.execute_input":"2022-07-15T07:37:04.829160Z","iopub.status.idle":"2022-07-15T07:37:04.842667Z","shell.execute_reply.started":"2022-07-15T07:37:04.829122Z","shell.execute_reply":"2022-07-15T07:37:04.841643Z"},"trusted":true},"execution_count":null,"outputs":[]},{"cell_type":"code","source":"class Model_hist(pl.LightningModule):\n    def __init__(self, learning_rate=1e-3):#, batch_size):\n        super().__init__()\n        self.model = HISTModel(d_feat=188, hidden_size=64, num_layers=2, dropout=0.3, base_model=\"GRU\")\n        self.learning_rate = learning_rate\n        self.train_amex_metric = AmexMetric()\n        self.val_amex_metric = AmexMetric()\n        self.loss_fn = nn.BCEWithLogitsLoss(reduction=\"mean\")\n        \n    def forward(self, x):\n        # in lightning, forward defines the prediction/inference actions\n        y_hat = self.model(x)\n        return y_hat\n\n    def training_step(self, batch, batch_idx):\n        # training_step defines the train loop. It is independent of forward\n        x, s2c, y = batch\n        x, s2c, y = x.float(), s2c.float(), y.float()\n        y_hat = self.model(x, s2c)\n        # loss function\n        loss = self.loss_fn(y_hat, y)\n        self.train_amex_metric.update(torch.sigmoid(y_hat), y)\n        self.log_dict({'train_loss': loss, 'train_amex_metric': self.train_amex_metric}, on_step=True, on_epoch=True, prog_bar=True, logger=True)\n        return {'loss': loss}\n    \n    def validation_step(self, batch, batch_idx):\n        # training_step defines the train loop. It is independent of forward\n        x, s2c, y = batch\n        x, s2c, y = x.float(), s2c.float(), y.float()\n        y_hat = self.model(x, s2c)\n        # loss function\n        loss = self.loss_fn(y_hat, y)\n        self.val_amex_metric.update(torch.sigmoid(y_hat), y)\n        self.log_dict({'val_loss': loss, 'val_amex_metric': self.val_amex_metric}, on_step=True, on_epoch=True, prog_bar=True, logger=True)\n        return {'loss': loss}       \n\n    def test_step(self, batch, batch_idx):\n        # training_step defines the train loop. It is independent of forward\n        x, s2c, y = batch\n        x, s2c, y = x.float(), s2c.float(), y.float()\n        y_hat = self.model(x, s2c)\n        # loss function\n        #loss = self.loss_fn(y_hats.squeeze(1), y_true)\n\n    def predict_step(self, batch, batch_idx):\n        # training_step defines the train loop. It is independent of forward\n        x, s2c, y = batch\n        x, s2c, y = x.float(), s2c.float(), y.float()\n        with torch.no_grad():\n            y_hat = self.model(x, s2c)#.squeeze(1)\n        return y_hat\n    \n    def configure_optimizers(self):\n        optimizer = torch.optim.AdamW(self.parameters(), lr=self.learning_rate)\n        #lr_scheduler = OneCycleLR(optimizer, max_lr=1e-3, epochs=25, steps_per_epoch=718)\n        lr_scheduler = CosineAnnealingWarmRestarts(optimizer, T_0=2, T_mult=3)\n        return [optimizer], [lr_scheduler]","metadata":{"execution":{"iopub.status.busy":"2022-07-15T07:37:19.298582Z","iopub.execute_input":"2022-07-15T07:37:19.299109Z","iopub.status.idle":"2022-07-15T07:37:19.315208Z","shell.execute_reply.started":"2022-07-15T07:37:19.299074Z","shell.execute_reply":"2022-07-15T07:37:19.314022Z"},"trusted":true},"execution_count":null,"outputs":[]},{"cell_type":"markdown","source":"# Training","metadata":{}},{"cell_type":"code","source":"prefix = 'hist'\n\nfor i in range(5):\n    dm = Dataset_pl(i)\n    model = Model_hist()#, argv['batch_size']) \n\n    #wandb_logger = WandbLogger(project=\"AMEX\")\n    callbacks=[ModelCheckpoint(dirpath='ckpt', \n                               monitor=\"val_amex_metric\", mode=\"max\")]\n\n#     trainer = pl.Trainer(gpus=[1], max_epochs=25, \n#                         logger=wandb_logger, callbacks=callbacks,\n#                         enable_progress_bar=True)\n    trainer = pl.Trainer(gpus=1, max_epochs=25, \n                         callbacks=callbacks,\n                         enable_progress_bar=True)\n    trainer.fit(model, datamodule=dm)\n\n    # get validation metrics\n    val = trainer.validate(model, datamodule=dm, ckpt_path='best')\n    val_amex_metric_epoch = val[0]['val_amex_metric_epoch']\n    \n    # get output\n    output = trainer.predict(model, datamodule=dm, ckpt_path='best')\n    output = torch.sigmoid(torch.cat(output))\n    # save result\n    df = pd.DataFrame(output)\n    df2 = pd.read_csv('submission.csv')\n    df2['prediction']=df.values\n    df2.to_csv(f'results/{prefix}_prediction_fold{i}.csv', index=False)\n    # save validation metrics to a txt file\n    with open(f'results/{prefix}_fold_{i}_{val_amex_metric_epoch}.txt', 'x') as f:\n        f.write(f\"fold {i} {val_amex_metric_epoch}\\n\")\n    del dm, trainer, model, output, df, df2\n    import gc\n    gc.collect()\n\n    print(f\"fold {i}\", val_amex_metric_epoch)","metadata":{"execution":{"iopub.status.busy":"2022-07-15T07:37:38.493261Z","iopub.execute_input":"2022-07-15T07:37:38.493750Z"},"trusted":true},"execution_count":null,"outputs":[]},{"cell_type":"code","source":"","metadata":{},"execution_count":null,"outputs":[]},{"cell_type":"code","source":"","metadata":{},"execution_count":null,"outputs":[]}]}