""" Preprocessor for Pump Sensor Data. Handles chronological sorting, sensor_15 drop, look-ahead window labeling, RECOVERING row exclusion, forward-fill + train-median imputation, temporal split, and StandardScaler fitting on train only. """ import os import sys import pandas as pd from sklearn.preprocessing import StandardScaler import joblib sys.path.insert(0, os.path.abspath(os.path.join(os.path.dirname(__file__), '..', 'utils'))) from preprocess_utils import logger, time_tracker @time_tracker def preprocess_pump(data_path: str, output_dir: str, lookahead_minutes: int = 60, split_ratio: float = 0.5) -> tuple: os.makedirs(output_dir, exist_ok=True) logger.info(f"Loading Pump Sensor dataset from {data_path}...") df = pd.read_csv(data_path) # Drop junk columns df = df.drop(columns=['Unnamed: 0'], errors='ignore') if 'sensor_15' in df.columns: df = df.drop(columns=['sensor_15']) logger.info("Dropped sensor_15 (100% null).") # Sort chronologically df['timestamp'] = pd.to_datetime(df['timestamp']) df = df.sort_values('timestamp').reset_index(drop=True) # Look-ahead window labeling broken_idx = df[df['machine_status'] == 'BROKEN'].index logger.info(f"Found {len(broken_idx)} BROKEN events.") df['failure_ahead'] = 0 for idx in broken_idx: start = max(0, idx - lookahead_minutes) df.loc[start:idx - 1, 'failure_ahead'] = 1 # Remove RECOVERING rows (leak source) n_recovering = (df['machine_status'] == 'RECOVERING').sum() logger.info(f"Removing {n_recovering} RECOVERING rows.") df = df[df['machine_status'] != 'RECOVERING'].copy() timestamps = df['timestamp'] targets = df['failure_ahead'] df_feat = df.drop(columns=['timestamp', 'machine_status', 'failure_ahead']) # Temporal split split_idx = int(len(df) * split_ratio) X_train, X_test = df_feat.iloc[:split_idx].copy(), df_feat.iloc[split_idx:].copy() y_train, y_test = targets.iloc[:split_idx], targets.iloc[split_idx:] # Impute: forward-fill then train-median fallback X_train, X_test = X_train.ffill(), X_test.ffill() medians = X_train.median() X_train, X_test = X_train.fillna(medians), X_test.fillna(medians) joblib.dump(medians, os.path.join(output_dir, 'pump_imputation_medians.pkl')) # Scale — fit on train ONLY scaler = StandardScaler() cols = list(X_train.columns) X_train[cols] = scaler.fit_transform(X_train[cols]) X_test[cols] = scaler.transform(X_test[cols]) joblib.dump(scaler, os.path.join(output_dir, 'pump_scaler.pkl')) # Save for tag, Xs, ys, ts in [('train', X_train, y_train, timestamps.iloc[:split_idx]), ('test', X_test, y_test, timestamps.iloc[split_idx:])]: out = Xs.copy() out['timestamp'], out['target'] = ts, ys out.to_csv(os.path.join(output_dir, f'pump_{tag}_processed.csv'), index=False) logger.info(f"Pump preprocessing complete. Saved to {output_dir}") return (os.path.join(output_dir, 'pump_train_processed.csv'), os.path.join(output_dir, 'pump_test_processed.csv'))