{"cells":[{"metadata":{"_uuid":"8f2839f25d086af736a60e9eeb907d3b93b6e0e5","_cell_guid":"b1076dfc-b9ad-4769-8c92-a6c4dae69d19","trusted":true},"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 in \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 \"../input/\" directory.\n# For example, running this (by clicking run or pressing Shift+Enter) will list the files in the input directory\n\nimport os\nprint(os.listdir(\"../input\"))\n\n# Any results you write to the current directory are saved as output.","execution_count":null,"outputs":[]},{"metadata":{"_cell_guid":"79c7e3d0-c299-4dcb-8224-4455121ee9b0","_uuid":"d629ff2d2480ee46fbb7e2d37f6b5fab8052498a","trusted":true},"cell_type":"code","source":"from typing import Tuple, Union\nimport pandas as pd\nimport numpy as np\nfrom tqdm import tqdm\nimport gc\nimport logging\n\n# logging.getLogger().setLevel(logging.DEBUG)\n\n\nclass FeatureExtraction:\n    \"\"\"\n    Class is responsible for creating features from the data.\n    \"\"\"\n\n    def __init__(self, ts_size: int, validation_fraction: float, segment_size: int = 150_000, ls_batch_size=None):\n        \"\"\"\n        Args:\n            ts_size: how many data points become one timestamp.\n            validation_fraction: fraction of data segments used for validation\n            segment_size: number of datapoints to be considered one segment. it is same as batch_size*ts_size.\n        \"\"\"\n        # Number of entries which make up one time stamp. note that features are learnt from this many datapoints\n        self._ts_size = ts_size\n        self._segment_size = segment_size\n        # just for test. In prod, 100*self._segment_size is used.\n        self._learn_scale_batch_size = ls_batch_size if ls_batch_size else 100 * self._segment_size\n\n        self._validation_fraction = validation_fraction\n\n        # to be set while fetching validation data/training data.\n        # number of training examples to be given to model\n        self.train_size = None\n        # number of validation examples to be given to model\n        self.validation_size = None\n        # number of datapoints in training data file.\n        self._raw_train_size = None\n\n        # Used for normalization. Contains the maximum absolute value for a feature.\n        self._scale_df = None\n\n        self._train_X_dfs = []\n        self._train_y_dfs = []\n\n        assert self._segment_size % self._ts_size == 0\n        assert self._validation_fraction >= 0\n\n    def get_y(self, df: pd.DataFrame) -> pd.Series:\n        df = df[['time_to_failure']].copy()\n        ts_count = df.shape[0] // self._ts_size\n        if df.shape[0] != ts_count * self._ts_size:\n            logging.warning('For y, Trimming last {} entries'.format(df.shape[0] - ts_count * self._ts_size))\n            df = df.iloc[:ts_count * self._ts_size]\n\n        df['ts'] = np.repeat(list(range(ts_count)), self._ts_size)\n        output_df = df.groupby('ts').last()\n        output_df.index.name = 'ts'\n        return output_df['time_to_failure']\n\n    @staticmethod\n    def group_location_filter(\n            df: pd.DataFrame,\n            grp_col: Union[str, int],\n            grp_size: int,\n            grp_start_index: int,\n            grp_end_index: int,\n    ) -> pd.core.groupby.DataFrameGroupBy:\n        \"\"\"\n        Within one, ts_size segment, we want to compute features on say first 10% of the data, or on some\n        contiguous segment. This function returns a group object which has exactly those entries.\n\n        Args:\n            grp_col: column name which is to be used to group the dataframe.\n            grp_size: how many entries are there in each group. It is assumed that it will be same for all groups.\n            grp_start_index: within a group, the index from which data needs to be considered.\n            grp_end_index: within a group, the index till which (excluding it) data needs to be considered.\n        \"\"\"\n        assert grp_start_index < df.shape[0] and grp_start_index >= 0\n        assert grp_end_index < df.shape[0]\n        idx = np.arange(0, df.shape[0])\n        df = df[(idx % grp_size >= grp_start_index) & (idx % grp_size < grp_end_index)]\n\n        return df.groupby(grp_col)\n\n    @staticmethod\n    def compute_features_on_group(grp: pd.core.groupby.DataFrameGroupBy, col_suffix: str):\n        mean_df = grp.mean().to_frame('mean_' + col_suffix)\n\n        if mean_df.empty:\n            return mean_df\n\n        std_df = grp.std().to_frame('std_' + col_suffix)\n        quantile_df = grp.quantile([0.05, 0.5, 0.95]).unstack()\n        quantile_df.columns = list(map(lambda x: 'Quantile-{}_{}'.format(x, col_suffix), quantile_df.columns))\n        max_df = grp.max().to_frame('max_' + col_suffix)\n        min_df = grp.min().to_frame('min_' + col_suffix)\n\n        output_df = pd.concat([mean_df, std_df, quantile_df, max_df, min_df], axis=1)\n        return output_df\n\n    def get_X(self, df: pd.DataFrame) -> pd.DataFrame:\n        \"\"\"\n        Index is #example_id,#timestamp_id. Columns are features.\n        \"\"\"\n        # TODO: this is where more features will come in.\n        df = df[['acoustic_data']].copy()\n        ts_count = df.shape[0] // self._ts_size\n\n        if df.shape[0] != ts_count * self._ts_size:\n            logging.warning('For X, Trimming last {} entries'.format(df.shape[0] - ts_count * self._ts_size))\n            df = df.iloc[:ts_count * self._ts_size]\n\n        df['ts'] = np.repeat(list(range(ts_count)), self._ts_size)\n\n        # Compute features on all datapoints within each ts_size chunk\n        grp = df.groupby('ts')\n        f_0_100_df = FeatureExtraction.compute_features_on_group(grp['acoustic_data'], '0->100')\n\n        # Compute features on first 25% of ts_size datapoints within each ts_size chunk\n        grp = FeatureExtraction.group_location_filter(df, 'ts', self._ts_size, 0, self._ts_size // 4)\n        f_0_25_df = FeatureExtraction.compute_features_on_group(grp['acoustic_data'], '0->25')\n\n        # Compute features on next 25% of ts_size datapoints within each ts_size chunk\n        grp = FeatureExtraction.group_location_filter(df, 'ts', self._ts_size, self._ts_size // 4, self._ts_size // 2)\n        f_25_50_df = FeatureExtraction.compute_features_on_group(grp['acoustic_data'], '25->50')\n\n        # Compute features on next 25% of ts_size datapoints within each ts_size chunk\n        grp = FeatureExtraction.group_location_filter(df, 'ts', self._ts_size, self._ts_size // 2,\n                                                      3 * self._ts_size // 4)\n        f_50_75_df = FeatureExtraction.compute_features_on_group(grp['acoustic_data'], '50->75')\n\n        # Compute features on next 25% of ts_size datapoints within each ts_size chunk\n        grp = FeatureExtraction.group_location_filter(df, 'ts', self._ts_size, 3 * self._ts_size // 4, self._ts_size)\n        f_75_100_df = FeatureExtraction.compute_features_on_group(grp['acoustic_data'], '75->100')\n\n        output_df = pd.concat([f_0_100_df, f_0_25_df, f_25_50_df, f_50_75_df, f_75_100_df], axis=1)\n        output_df.columns.name = 'features'\n\n        if self._scale_df is not None:\n            output_df = output_df / self._scale_df\n\n        gc.collect()\n        return output_df\n\n    def learn_scale_and_save_train_df(self, raw_data_df):\n        \"\"\"\n        Learns scale for normalization from training data. it does not touch validation data.\n        \"\"\"\n        logging.info('Scale about to be learnt')\n        self._train_X_dfs = []\n        self._train_y_dfs = []\n\n        train_X_dfs = []\n        train_y_dfs = []\n\n        self._scale_df = scale_df = None\n        gen = self.get_X_y_generator(raw_data_df, 0, batch_size=self._learn_scale_batch_size)\n\n        for X_df, y_df in gen:\n            gc.collect()\n            train_X_dfs.append(X_df)\n            train_y_dfs.append(y_df)\n\n            max_df = X_df.abs().max()\n            if scale_df is None:\n                scale_df = max_df\n            else:\n                scale_df = pd.concat([scale_df, max_df], axis=1).max(axis=1)\n\n        self._scale_df = scale_df\n\n        logging.info('Scale learnt')\n        self._train_X_dfs = [df / self._scale_df for df in train_X_dfs]\n        self._train_y_dfs = train_y_dfs\n\n    def get_X_y(self, df: pd.DataFrame) -> Tuple[pd.DataFrame, pd.DataFrame]:\n        X_df = self.get_X(df)\n        y_df = self.get_y(df)\n        return (X_df, y_df)\n\n    def _set_train_validation_size(self):\n        # We ensure that last few segments are not used in training data. We use it for validation.\n        num_segments = self._raw_train_size // self._segment_size\n\n        validation_segments = int(self._validation_fraction * num_segments)\n        train_segments = num_segments - validation_segments\n\n        self.train_size = int(train_segments * self._segment_size / self._ts_size)\n        self.validation_size = int(validation_segments * self._segment_size / self._ts_size)\n        print('Validation Size', self.validation_size)\n        print('Train Size', self.train_size)\n\n    def get_validation_X_y(self, raw_data_df):\n        logging.info('Validation data requested.')\n        if self._validation_fraction == 0:\n            return (pd.DataFrame(), pd.Series())\n\n        if self._raw_train_size is None:\n            self._raw_train_size = raw_data_df.shape[0]\n            self._set_train_validation_size()\n\n        val_Xs = []\n        val_ys = []\n        for start_index in range(self.train_size * self._ts_size, raw_data_df.shape[0], self._segment_size):\n            df = raw_data_df.iloc[start_index:start_index + self._segment_size]\n            val_X_df, val_y_df = self.get_X_y(df)\n            val_Xs.append(val_X_df)\n            val_ys.append(val_y_df)\n\n        val_X_df = pd.concat(val_Xs)\n        val_y_df = pd.concat(val_ys)\n\n        logging.info('Validation data returned.')\n        return (val_X_df, val_y_df)\n\n    def get_X_y_generator_fast(self, padding_row_count, debug_mode=False):\n        assert self._learn_scale_batch_size % self._segment_size == 0\n        assert padding_row_count % self._ts_size == 0\n        assert len(self._train_X_dfs) > 0 or self._raw_train_size < self._ts_size\n\n        padding_row_count = padding_row_count // self._ts_size\n        segment_size_row_count = self._segment_size // self._ts_size\n\n        while True:\n            prev_X_df = None\n            prev_y_df = None\n            for X_chunk_df, y_chunk_df in zip(self._train_X_dfs, self._train_y_dfs):\n                for start_index in range(0, X_chunk_df.shape[0], segment_size_row_count):\n                    padding_start_index = start_index - padding_row_count\n\n                    if padding_start_index >= 0:\n                        X_df = X_chunk_df.iloc[padding_start_index:start_index + segment_size_row_count]\n                        y_df = y_chunk_df.iloc[padding_start_index:start_index + segment_size_row_count]\n                        prev_X_df = X_df\n                        prev_y_df = y_df\n                        yield (X_df, y_df)\n                    else:\n                        X_df = X_chunk_df.iloc[:start_index + segment_size_row_count]\n                        y_df = y_chunk_df.iloc[:start_index + segment_size_row_count]\n\n                        if prev_X_df is not None:\n                            X_df = pd.concat([prev_X_df.iloc[padding_start_index:], X_df])\n                            y_df = pd.concat([prev_y_df.iloc[padding_start_index:], y_df])\n\n                        prev_X_df = X_df\n                        prev_y_df = y_df\n                        yield (X_df, y_df)\n\n            if debug_mode:\n                break\n\n    def get_X_y_generator(self, raw_data_df, padding_row_count: int,\n                          batch_size=None) -> Tuple[pd.DataFrame, pd.DataFrame]:\n\n        assert padding_row_count % self._ts_size == 0\n        if batch_size is None:\n            batch_size = self._segment_size\n\n        if self._raw_train_size is None:\n            self._raw_train_size = raw_data_df.shape[0]\n            self._set_train_validation_size()\n\n        logging.info('Training X,y generator starting from beginning')\n        next_first_index = 0\n        for start_index in range(0, self.train_size * self._ts_size, batch_size):\n            df = raw_data_df.iloc[max(0, start_index - padding_row_count):start_index + batch_size]\n\n            X_df, y_df = self.get_X_y(df)\n            X_df.index += next_first_index\n            y_df.index += next_first_index\n            gc.collect()\n\n            # if it padding is non zero, few entries from last segment will come in next segment\n            padded_first_entry_index = (1 + padding_row_count // self._ts_size)\n            # the if-else is a corner case. If padding is more than segment_size then this will happen.\n            if padding_row_count == 0:\n                next_first_index = y_df.index[-1] + 1\n            elif padded_first_entry_index <= y_df.shape[0]:\n                next_first_index = y_df.index[-1 * padded_first_entry_index + 1]\n            else:\n                next_first_index = 0\n\n            yield (X_df, y_df)\n\n\nclass Data:\n    \"\"\"\n    This class uses FeatureExtraction class and creates a time sequence data.\n    \"\"\"\n\n    def __init__(\n            self,\n            ts_window: int,\n            ts_size: int,\n            train_fname: str,\n            normalize: bool = True,\n            validation_fraction: float = 0.2,\n            segment_size: int = 150_000,\n            ls_batch_size: int = None,\n    ):\n        self._ts_window = ts_window\n        self._ts_size = ts_size\n        self._validation_fraction = validation_fraction\n        self._segment_size = segment_size\n        self._train_fname = train_fname\n        self.raw_data_df = pd.read_csv(\n            train_fname,\n            dtype={'acoustic_data': np.int16,\n                   'time_to_failure': np.float32},\n        )\n\n        self._feature_extractor = FeatureExtraction(\n            self._ts_size,\n            self._validation_fraction,\n            segment_size=segment_size,\n            ls_batch_size=ls_batch_size,\n        )\n\n        self._normalize = normalize\n\n        if self._normalize:\n            # Learn scale.\n            self._feature_extractor.learn_scale_and_save_train_df(self.raw_data_df)\n            print('[Data] Scale learnt')\n\n        print('[Data] Fetching Validation data')\n        self.val_X, self.val_y = self.get_validation_X_y()\n        print('[Data] Validation data fetched')\n\n    def get_window_X(self, X_df: pd.DataFrame) -> np.array:\n        row_count = X_df.shape[0] - self._ts_window + 1\n\n        if row_count <= 0:\n            return None\n\n        X = np.zeros((row_count, self._ts_window, X_df.shape[1]))\n        for i in range(self._ts_window, X_df.shape[0] + 1):\n            X[i - self._ts_window] = X_df.values[i - self._ts_window:i, :]\n        return X\n\n    def training_size(self):\n        \"\"\"\n        Returns number of examples to be used in training.\n        \"\"\"\n        return self._feature_extractor.train_size\n\n    def batch_size(self):\n        # 150_000\n        return self._segment_size // self._ts_size\n\n    def get_window_y(self, y_df: pd.Series) -> np.array:\n        return y_df.values[self._ts_window - 1:]\n\n    def get_window_X_y(self, X_df, y_df) -> Tuple[np.array, np.array]:\n        X = self.get_window_X(X_df)\n        y = self.get_window_y(y_df)\n\n        if X is None:\n            return (None, None)\n\n        return (X, y)\n\n    def get_validation_X_y(self) -> Tuple[np.array, np.array]:\n        \"\"\"\n        Returns last few segments of training data for validation. Note that this is not used in\n        training, ie, this is not returned from get_X_y_generator()\n        \"\"\"\n\n        X_df, y_df = self._feature_extractor.get_validation_X_y(self.raw_data_df)\n        X, y = self.get_window_X_y(X_df, y_df)\n        return (X, y)\n\n    def get_test_X(self, df) -> np.array:\n        \"\"\"\n        For Test data, it fetches data from file and returns the X with shape\n             (#examples, self._ts_window, feature_count)\n        \"\"\"\n        gc.collect()\n        X_df = self._feature_extractor.get_X(df)\n        return self.get_window_X(X_df)\n\n    def get_X_y_generator(self, debug_mode: bool = False) -> Tuple[np.array, np.array]:\n        # We need self._ts_window -1 rows at beginning to cater to starting data points in a chunk.\n        padding = self._ts_size * (self._ts_window - 1)\n        # gen = self._feature_extractor.get_X_y_generator(self.raw_data_df, padding, test_mode=test_mode)\n        gen = self._feature_extractor.get_X_y_generator_fast(padding, debug_mode=debug_mode)\n\n        for X_df, y_df in gen:\n            X, y = self.get_window_X_y(X_df, y_df)\n            yield (X, y)\n\n\n# if __name__ == '__main__':\n#     ts_window = 100\n#     ts_size = 1000\n#     d = Data(ts_window, ts_size, 'train.csv')\n#     gen = d.get_X_y_generator(test_mode=True)\n#     for X, y in gen:\n#         print('Shape of X', X.shape)\n#         print('Shape of y', y.shape)\n","execution_count":1,"outputs":[]},{"metadata":{"trusted":true},"cell_type":"code","source":"import pickle \nts_window = 150\nts_size = 1000\nd = Data(ts_window, ts_size, '../input/train.csv')\nwith open('train_data.pkl', 'wb') as f:\n    pickle.dump(d, f)","execution_count":null,"outputs":[]}],"metadata":{"kernelspec":{"display_name":"Python 3","language":"python","name":"python3"},"language_info":{"name":"python","version":"3.6.4","mimetype":"text/x-python","codemirror_mode":{"name":"ipython","version":3},"pygments_lexer":"ipython3","nbconvert_exporter":"python","file_extension":".py"}},"nbformat":4,"nbformat_minor":1}