""" Preprocessor for NASA CMAPSS Turbofan Engine Dataset. Computes RUL labels, applies MinMaxScaler (fit on train), and generates sliding-window numpy sequences for LSTM input. """ import os import sys import numpy as np import pandas as pd from sklearn.preprocessing import MinMaxScaler import joblib sys.path.insert(0, os.path.abspath(os.path.join(os.path.dirname(__file__), '..', 'utils'))) from preprocess_utils import logger, time_tracker def _generate_sequences(df, window_size, feature_cols, sub_dataset): """Build sliding-window arrays grouped by engine (no cross-engine leakage).""" X_seq, y_seq = [], [] for eid, grp in df.groupby('engine_id'): if len(grp) < window_size: logger.warning(f"Engine {eid} ({sub_dataset}): {len(grp)} cycles < window {window_size}, skipped.") continue feats = grp[feature_cols].values ruls = grp['RUL'].values for i in range(len(grp) - window_size + 1): X_seq.append(feats[i:i + window_size]) y_seq.append(ruls[i + window_size - 1]) return np.array(X_seq), np.array(y_seq) @time_tracker def preprocess_cmapss(data_dir: str, output_dir: str, sub_dataset: str = 'FD001', window_size: int = 30) -> tuple: os.makedirs(output_dir, exist_ok=True) logger.info(f"Processing CMAPSS {sub_dataset}...") idx_cols = ['engine_id', 'cycle'] set_cols = ['setting1', 'setting2', 'setting3'] sen_cols = [f'sensor{i}' for i in range(1, 22)] col_names = idx_cols + set_cols + sen_cols paths = {k: os.path.join(data_dir, f'{k}_{sub_dataset}.txt') for k in ['train', 'test', 'RUL']} for p in paths.values(): if not os.path.exists(p): raise FileNotFoundError(f"Missing: {p}") df_tr = pd.read_csv(paths['train'], sep=r'\s+', header=None, names=col_names) df_te = pd.read_csv(paths['test'], sep=r'\s+', header=None, names=col_names) df_rul = pd.read_csv(paths['RUL'], sep=r'\s+', header=None, names=['RUL_gt']) # RUL for train: max_cycle - current_cycle max_cyc = df_tr.groupby('engine_id')['cycle'].max().reset_index() max_cyc.columns = ['engine_id', 'max_cycle'] df_tr = df_tr.merge(max_cyc, on='engine_id') df_tr['RUL'] = df_tr['max_cycle'] - df_tr['cycle'] df_tr.drop(columns='max_cycle', inplace=True) # RUL for test: ground_truth + max_cycle - cycle df_rul['engine_id'] = df_rul.index + 1 max_cyc_te = df_te.groupby('engine_id')['cycle'].max().reset_index() max_cyc_te.columns = ['engine_id', 'max_cycle'] df_te = df_te.merge(max_cyc_te, on='engine_id').merge(df_rul, on='engine_id') df_te['RUL'] = df_te['RUL_gt'] + df_te['max_cycle'] - df_te['cycle'] df_te.drop(columns=['max_cycle', 'RUL_gt'], inplace=True) # Scale — fit on train only scale_cols = set_cols + sen_cols scaler = MinMaxScaler() df_tr[scale_cols] = scaler.fit_transform(df_tr[scale_cols]) df_te[scale_cols] = scaler.transform(df_te[scale_cols]) joblib.dump(scaler, os.path.join(output_dir, f'cmapss_{sub_dataset}_scaler.pkl')) # Generate sequences X_tr, y_tr = _generate_sequences(df_tr, window_size, scale_cols, sub_dataset) X_te, y_te = _generate_sequences(df_te, window_size, scale_cols, sub_dataset) out_paths = [] for name, arr in [('X_train', X_tr), ('y_train', y_tr), ('X_test', X_te), ('y_test', y_te)]: p = os.path.join(output_dir, f'cmapss_{sub_dataset}_{name}.npy') np.save(p, arr) out_paths.append(p) logger.info(f"CMAPSS {sub_dataset} done. Train: X={X_tr.shape}, Test: X={X_te.shape}") return tuple(out_paths)