""" Preprocessor for Microsoft Azure Predictive Maintenance Dataset. Merges telemetry, errors, maintenance, failures, and machine metadata. Engineers 24h rolling features, creates look-ahead failure labels, and applies temporal split + StandardScaler (train-only fit). """ import os import sys import pandas as pd from sklearn.preprocessing import LabelEncoder, 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_azure(data_dir: str, output_dir: str, split_date: str = '2015-09-01') -> tuple: os.makedirs(output_dir, exist_ok=True) logger.info("Loading Azure PdM files...") tele = pd.read_csv(os.path.join(data_dir, 'PdM_telemetry.csv')) errs = pd.read_csv(os.path.join(data_dir, 'PdM_errors.csv')) mnt = pd.read_csv(os.path.join(data_dir, 'PdM_maint.csv')) fail = pd.read_csv(os.path.join(data_dir, 'PdM_failures.csv')) mach = pd.read_csv(os.path.join(data_dir, 'PdM_machines.csv')) for d in [tele, errs, mnt, fail]: d['datetime'] = pd.to_datetime(d['datetime']) # --- 1. Rolling 24h telemetry features --- logger.info("Computing 24h rolling telemetry features...") sensors = ['volt', 'rotate', 'pressure', 'vibration'] parts = [] for mid, grp in tele.groupby('machineID'): grp = grp.sort_values('datetime').set_index('datetime')[sensors] roll_mean = grp.rolling(24, min_periods=1).mean() roll_std = grp.rolling(24, min_periods=1).std().fillna(0) roll_mean.columns = [f'{c}_rollmean_24h' for c in sensors] roll_std.columns = [f'{c}_rollstd_24h' for c in sensors] combined = pd.concat([roll_mean, roll_std], axis=1).reset_index() combined['machineID'] = mid parts.append(combined) feat = pd.concat(parts) # --- 2. Error log counts --- logger.info("Merging error counts...") err_oh = pd.get_dummies(errs, columns=['errorID']).groupby(['machineID', 'datetime']).sum().reset_index() feat = feat.merge(err_oh, on=['machineID', 'datetime'], how='left') err_cols = [c for c in feat.columns if 'errorID_' in c] feat[err_cols] = feat[err_cols].fillna(0) # --- 3. Maintenance log counts --- logger.info("Merging maintenance records...") mnt_oh = pd.get_dummies(mnt, columns=['comp']).groupby(['machineID', 'datetime']).sum().reset_index() feat = feat.merge(mnt_oh, on=['machineID', 'datetime'], how='left') mnt_cols = [c for c in feat.columns if 'comp_' in c] feat[mnt_cols] = feat[mnt_cols].fillna(0) # --- 4. Machine metadata --- logger.info("Merging machine metadata...") mach['model_encoded'] = LabelEncoder().fit_transform(mach['model']) feat = feat.merge(mach[['machineID', 'age', 'model_encoded']], on='machineID', how='left') # --- 5. 24h look-ahead failure target --- logger.info("Creating 24h look-ahead failure labels...") fail['failure_target'] = 1 feat = feat.merge(fail[['machineID', 'datetime', 'failure_target']], on=['machineID', 'datetime'], how='left') feat['failure_target'] = feat['failure_target'].fillna(0) feat['target'] = 0 for _, grp in feat.groupby('machineID'): for idx in grp[grp['failure_target'] == 1].index: start = max(grp.index[0], idx - 24) feat.loc[start:idx, 'target'] = 1 feat.drop(columns='failure_target', inplace=True) # --- 6. Temporal split + scale --- split_dt = pd.to_datetime(split_date) train_df = feat[feat['datetime'] < split_dt].copy() test_df = feat[feat['datetime'] >= split_dt].copy() num_cols = [f'{s}_{t}_24h' for s in sensors for t in ['rollmean', 'rollstd']] + ['age'] scaler = StandardScaler() train_df[num_cols] = scaler.fit_transform(train_df[num_cols]) test_df[num_cols] = scaler.transform(test_df[num_cols]) joblib.dump(scaler, os.path.join(output_dir, 'azure_scaler.pkl')) tr_path = os.path.join(output_dir, 'azure_train_processed.csv') te_path = os.path.join(output_dir, 'azure_test_processed.csv') train_df.to_csv(tr_path, index=False) test_df.to_csv(te_path, index=False) logger.info(f"Azure PdM done. Train: {train_df.shape} ({train_df['target'].mean()*100:.2f}% fail), " f"Test: {test_df.shape} ({test_df['target'].mean()*100:.2f}% fail)") return tr_path, te_path