File size: 4,469 Bytes
baf834b | 1 2 3 4 5 6 7 8 9 10 11 12 13 14 15 16 17 18 19 20 21 22 23 24 25 26 27 28 29 30 31 32 33 34 35 36 37 38 39 40 41 42 43 44 45 46 47 48 49 50 51 52 53 54 55 56 57 58 59 60 61 62 63 64 65 66 67 68 69 70 71 72 73 74 75 76 77 78 79 80 81 82 83 84 85 86 87 88 89 90 91 92 93 94 95 96 97 98 99 100 | """
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
|