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'))
|