SyedaArisha's picture
Upload folder using huggingface_hub
baf834b verified
Raw History Blame Contribute Delete
3.16 kB
"""
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'))