File size: 3,161 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
"""
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'))