""" Data preprocessing and feature engineering utilities for job failure prediction. """ import pandas as pd import numpy as np from typing import Dict, List, Optional import json def parse_duration(x) -> float: """Parse duration string (HH:MM:SS) or numeric to seconds.""" if pd.isna(x): return 0.0 try: if isinstance(x, str): parts = x.split(':') if len(parts) == 3: h, m, s = map(int, parts) return h * 3600 + m * 60 + s return float(x) return float(x) except (ValueError, AttributeError): return 0.0 def engineer_features(df: pd.DataFrame) -> pd.DataFrame: """ Engineer features from raw job data. Args: df: DataFrame with columns: zone, job_nm, tasksgroup_nm, round_time, job_start_time, job_end_time, duration, status, err_msg, zeppelin, ictrl_dt, start_ictrl_dt, end_ictrl_dt Returns: DataFrame with engineered features """ df = df.copy() # Parse timestamps df['job_start_time'] = pd.to_datetime(df['job_start_time'], errors='coerce') df['job_end_time'] = pd.to_datetime(df['job_end_time'], errors='coerce') # Parse duration to seconds if 'duration' in df.columns: df['duration_sec'] = df['duration'].apply(parse_duration) else: df['duration_sec'] = 0.0 df['duration_sec'] = df['duration_sec'].fillna(0.0) # Ground truth label # Handle various success statuses: 'SUCCESS', 'SUCCEED', 'SUCCEEDED' # Everything else (FAILED, ABORT-AUTO, RUNNING, etc.) is considered a failure status_upper = df['status'].fillna('').str.upper() success_statuses = ['SUCCESS', 'SUCCEED', 'SUCCEEDED'] df['is_failed'] = (~status_upper.isin(success_statuses)).astype(int) # Time features df['run_hour'] = df['job_start_time'].dt.hour.fillna(0).astype(int) df['run_dow'] = df['job_start_time'].dt.dayofweek.fillna(0).astype(int) # 0=Mon, 6=Sun df['is_weekend'] = (df['run_dow'] >= 5).astype(int) # Cyclical encoding for hour df['hour_sin'] = np.sin(2 * np.pi * df['run_hour'] / 24) df['hour_cos'] = np.cos(2 * np.pi * df['run_hour'] / 24) # Error message features df['err_msg_len'] = df['err_msg'].fillna('').str.len() df['has_err_msg'] = (df['err_msg_len'] > 0).astype(int) # Zeppelin flag df['is_zeppelin'] = df['zeppelin'].notna().astype(int) # Job-level rolling statistics (per job_nm) df = df.sort_values(['job_nm', 'job_start_time']).reset_index(drop=True) df['failure_rate_7'] = df.groupby('job_nm')['is_failed'].transform( lambda s: s.rolling(7, min_periods=1).mean() ) df['avg_duration_7'] = df.groupby('job_nm')['duration_sec'].transform( lambda s: s.rolling(7, min_periods=1).mean() ) # Duration z-score (relative to rolling average) df['duration_zscore'] = ( (df['duration_sec'] - df['avg_duration_7']) / df['avg_duration_7'].replace(0, 1) ) df['duration_zscore'] = df['duration_zscore'].fillna(0.0) return df def get_feature_columns() -> Dict[str, List[str]]: """Return feature column definitions.""" return { 'numeric': [ 'duration_sec', 'duration_zscore', 'avg_duration_7', 'failure_rate_7', 'err_msg_len', 'hour_sin', 'hour_cos' ], 'categorical': [ 'job_nm', 'tasksgroup_nm', 'zone', 'is_zeppelin', 'is_weekend' ], 'anomaly_numeric': [ 'duration_sec', 'duration_zscore', 'avg_duration_7', 'failure_rate_7', 'err_msg_len', 'hour_sin', 'hour_cos' ] } def save_feature_schema(output_path: str = 'models/feature_schema.json'): """Save feature schema to JSON file.""" import os os.makedirs(os.path.dirname(output_path), exist_ok=True) schema = { 'feature_columns': get_feature_columns(), 'required_fields': [ 'zone', 'job_nm', 'tasksgroup_nm', 'job_start_time', 'duration', 'status', 'err_msg', 'zeppelin' ] } with open(output_path, 'w') as f: json.dump(schema, f, indent=2) return schema def align_schema_df(df: pd.DataFrame, schema_path: str = 'models/feature_schema.json') -> pd.DataFrame: """ Align DataFrame to feature schema, filling missing columns with defaults. Args: df: Input DataFrame schema_path: Path to feature schema JSON Returns: Aligned DataFrame """ try: with open(schema_path, 'r') as f: schema = json.load(f) except FileNotFoundError: # If schema doesn't exist, engineer features and create it df = engineer_features(df) schema = save_feature_schema(schema_path) # Ensure all required numeric and categorical columns exist all_features = schema['feature_columns']['numeric'] + schema['feature_columns']['categorical'] for col in all_features: if col not in df.columns: if col in schema['feature_columns']['numeric']: df[col] = 0.0 else: df[col] = '' return df