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