Spaces:
Sleeping
Sleeping
| """ | |
| 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 | |