{"metadata":{"kernelspec":{"language":"python","display_name":"Python 3","name":"python3"},"language_info":{"name":"python","version":"3.11.11","mimetype":"text/x-python","codemirror_mode":{"name":"ipython","version":3},"pygments_lexer":"ipython3","nbconvert_exporter":"python","file_extension":".py"},"kaggle":{"accelerator":"none","dataSources":[{"sourceId":96164,"databundleVersionId":11418275,"sourceType":"competition"}],"dockerImageVersionId":31041,"isInternetEnabled":true,"language":"python","sourceType":"notebook","isGpuEnabled":false}},"nbformat_minor":4,"nbformat":4,"cells":[{"cell_type":"code","source":"# This Python 3 environment comes with many helpful analytics libraries installed\n# It is defined by the kaggle/python Docker image: https://github.com/kaggle/docker-python\n# For example, here's several helpful packages to load\n\nimport numpy as np # linear algebra\nimport pandas as pd # data processing, CSV file I/O (e.g. pd.read_csv)\n\n# Input data files are available in the read-only \"../input/\" directory\n# For example, running this (by clicking run or pressing Shift+Enter) will list all files under the input directory\n\nimport os\nfor dirname, _, filenames in os.walk('/kaggle/input'):\n    for filename in filenames:\n        print(os.path.join(dirname, filename))\n\n# You can write up to 20GB to the current directory (/kaggle/working/) that gets preserved as output when you create a version using \"Save & Run All\" \n# You can also write temporary files to /kaggle/temp/, but they won't be saved outside of the current session","metadata":{"_uuid":"8f2839f25d086af736a60e9eeb907d3b93b6e0e5","_cell_guid":"b1076dfc-b9ad-4769-8c92-a6c4dae69d19","trusted":true,"execution":{"iopub.status.busy":"2025-06-01T04:25:41.408260Z","iopub.execute_input":"2025-06-01T04:25:41.408481Z","iopub.status.idle":"2025-06-01T04:25:41.667054Z","shell.execute_reply.started":"2025-06-01T04:25:41.408464Z","shell.execute_reply":"2025-06-01T04:25:41.666353Z"}},"outputs":[],"execution_count":null},{"cell_type":"code","source":"#!/usr/bin/env python\n# -*- coding: utf-8 -*-\n\n\"\"\"\nDRW Crypto Market Prediction - Extended Runtime Pipeline with Comprehensive Checkpointing\nThis implementation extends runtime significantly and provides full pipeline state persistence\nallowing resumption from any stage of execution.\n\"\"\"\n\n# Install required packages\n!pip install flaml -q\n!pip install lightgbm -q\n!pip install xgboost -q\n!pip install catboost -q\n!pip install psutil -q\n\nimport pandas as pd\nimport numpy as np\nfrom flaml import AutoML\nimport lightgbm as lgb\nfrom sklearn.model_selection import KFold\nfrom scipy.stats import pearsonr\nimport warnings\nimport json\nfrom datetime import datetime, timedelta\nimport gc\nimport itertools\nimport os\nimport psutil\nimport pickle\nimport time\nimport hashlib\n\nwarnings.filterwarnings('ignore')\n\nclass CFG:\n    \"\"\"Configuration settings for extended runtime pipeline\"\"\"\n    # File paths\n    train_path = \"/kaggle/input/drw-crypto-market-prediction/train.parquet\"\n    test_path = \"/kaggle/input/drw-crypto-market-prediction/test.parquet\"\n    sample_sub_path = \"/kaggle/input/drw-crypto-market-prediction/sample_submission.csv\"\n    \n    # Extended checkpoint paths\n    checkpoint_dir = \"./checkpoints\"\n    checkpoint_version = \"v2\"\n    checkpoint_manifest = \"checkpoint_manifest.json\"\n    checkpoint_features = \"checkpoint_features.pkl\"\n    checkpoint_history = \"checkpoint_history.json\"\n    checkpoint_data = \"checkpoint_data.pkl\"\n    checkpoint_model = \"checkpoint_model.pkl\"\n    checkpoint_pipeline = \"checkpoint_pipeline.json\"\n    checkpoint_results = \"checkpoint_results.pkl\"\n    \n    # Model parameters\n    n_folds = 10  # Increased from 5\n    random_seed = 42\n    \n    # Significantly extended time budgets\n    feature_evaluation_time = 500  # 10 minutes per feature evaluation (was 300)\n    final_model_time = 14400*.5*.5  # 4 hours for final model (was 3600)\n    max_runtime_hours = 11.5  # Maximum total runtime in hours\n    \n    # Expanded feature engineering parameters\n    new_feature_depth = 6  # Increased from 6\n    feature_width = 60  # Increased from 40\n    min_improvement = 0.00001  # Reduced from 0.00005\n    max_base_features = 150  # Increased from 100\n    \n    # Enhanced transformations\n    transformations = ['log', 'sqrt', 'square', 'reciprocal', 'exp', 'abs', 'tanh', 'sigmoid']\n    transformation_weight = 0.35\n    \n    # Additional operations for combinations\n    advanced_operations = ['max', 'min', 'mean', 'std', 'weighted_mean']\n    \n    # Performance optimization\n    parallel_evaluation = True\n    memory_optimization_interval = 50  # Run memory optimization every N features\n    \n    # Checkpoint intervals\n    checkpoint_interval_minutes = 30  # Save checkpoint every 30 minutes\n    feature_checkpoint_interval = 1  # Save after every 5 features added\n\nclass PipelineState:\n    \"\"\"Manages pipeline execution state\"\"\"\n    def __init__(self):\n        self.start_time = None\n        self.stage = None\n        self.substage = None\n        self.progress = {}\n        self.results = {}\n        self.errors = []\n        self.runtime_stats = {}\n        \n    def to_dict(self):\n        return {\n            'start_time': self.start_time.isoformat() if self.start_time else None,\n            'stage': self.stage,\n            'substage': self.substage,\n            'progress': self.progress,\n            'results': self.results,\n            'errors': self.errors,\n            'runtime_stats': self.runtime_stats\n        }\n    \n    @classmethod\n    def from_dict(cls, data):\n        state = cls()\n        state.start_time = datetime.fromisoformat(data['start_time']) if data.get('start_time') else None\n        state.stage = data.get('stage')\n        state.substage = data.get('substage')\n        state.progress = data.get('progress', {})\n        state.results = data.get('results', {})\n        state.errors = data.get('errors', [])\n        state.runtime_stats = data.get('runtime_stats', {})\n        return state\n\ndef ensure_checkpoint_dir():\n    \"\"\"Ensure checkpoint directory exists\"\"\"\n    if not os.path.exists(CFG.checkpoint_dir):\n        os.makedirs(CFG.checkpoint_dir)\n        print(f\"Created checkpoint directory: {CFG.checkpoint_dir}\")\n\ndef get_checkpoint_path(filename):\n    \"\"\"Get full checkpoint path with version\"\"\"\n    return os.path.join(CFG.checkpoint_dir, f\"{CFG.checkpoint_version}_{filename}\")\n\ndef save_checkpoint(checkpoint_type, data, use_versioning=True):\n    \"\"\"Save checkpoint data with optional versioning\"\"\"\n    ensure_checkpoint_dir()\n    \n    if use_versioning:\n        filepath = get_checkpoint_path(checkpoint_type)\n    else:\n        filepath = os.path.join(CFG.checkpoint_dir, checkpoint_type)\n    \n    try:\n        if checkpoint_type.endswith('.json'):\n            with open(filepath, 'w') as f:\n                json.dump(data, f, indent=2, default=str)\n        else:\n            with open(filepath, 'wb') as f:\n                pickle.dump(data, f)\n        print(f\"Checkpoint saved: {filepath}\")\n        return True\n    except Exception as e:\n        print(f\"Error saving checkpoint {checkpoint_type}: {e}\")\n        return False\n\ndef load_checkpoint(checkpoint_type, use_versioning=True):\n    \"\"\"Load checkpoint data if exists\"\"\"\n    if use_versioning:\n        filepath = get_checkpoint_path(checkpoint_type)\n    else:\n        filepath = os.path.join(CFG.checkpoint_dir, checkpoint_type)\n    \n    if not os.path.exists(filepath):\n        return None\n    \n    try:\n        if checkpoint_type.endswith('.json'):\n            with open(filepath, 'r') as f:\n                data = json.load(f)\n        else:\n            with open(filepath, 'rb') as f:\n                data = pickle.load(f)\n        print(f\"Checkpoint loaded: {filepath}\")\n        return data\n    except Exception as e:\n        print(f\"Error loading checkpoint {checkpoint_type}: {e}\")\n        return None\n\ndef save_pipeline_state(state):\n    \"\"\"Save complete pipeline state\"\"\"\n    manifest = {\n        'version': CFG.checkpoint_version,\n        'timestamp': datetime.now().isoformat(),\n        'state': state.to_dict(),\n        'config': {\n            'feature_depth': CFG.new_feature_depth,\n            'feature_width': CFG.feature_width,\n            'min_improvement': CFG.min_improvement,\n            'max_runtime_hours': CFG.max_runtime_hours\n        }\n    }\n    return save_checkpoint(CFG.checkpoint_manifest, manifest, use_versioning=False)\n\ndef load_pipeline_state():\n    \"\"\"Load complete pipeline state\"\"\"\n    manifest = load_checkpoint(CFG.checkpoint_manifest, use_versioning=False)\n    if manifest:\n        return PipelineState.from_dict(manifest['state'])\n    return PipelineState()\n\ndef get_memory_usage():\n    \"\"\"Get current memory usage in GB\"\"\"\n    process = psutil.Process(os.getpid())\n    return process.memory_info().rss / 1024 / 1024 / 1024\n\ndef check_runtime_limit(start_time):\n    \"\"\"Check if runtime limit has been reached\"\"\"\n    elapsed = (datetime.now() - start_time).total_seconds() / 3600\n    remaining = CFG.max_runtime_hours - elapsed\n    \n    if remaining <= 0:\n        print(f\"\\nRuntime limit reached ({CFG.max_runtime_hours} hours)\")\n        return True, 0\n    \n    print(f\"Runtime: {elapsed:.1f}/{CFG.max_runtime_hours} hours ({remaining:.1f} hours remaining)\")\n    return False, remaining\n\ndef safe_divide(numerator, denominator, fill_value=0):\n    \"\"\"Safe division handling zero denominators\"\"\"\n    with np.errstate(divide='ignore', invalid='ignore'):\n        result = np.where(\n            np.abs(denominator) > 1e-10,\n            numerator / denominator,\n            fill_value\n        )\n    return np.nan_to_num(result, nan=fill_value, posinf=fill_value, neginf=fill_value)\n\ndef reduce_memory_usage(df, name):\n    \"\"\"Optimize dataframe memory usage\"\"\"\n    start_mem = df.memory_usage().sum() / 1024**2\n    print(f'Memory usage of {name}: {start_mem:.2f} MB')\n    \n    for col in df.columns:\n        col_type = df[col].dtype\n        if col_type != object:\n            c_min = df[col].min()\n            c_max = df[col].max()\n            \n            if str(col_type)[:3] == 'int':\n                if c_min > np.iinfo(np.int8).min and c_max < np.iinfo(np.int8).max:\n                    df[col] = df[col].astype(np.int8)\n                elif c_min > np.iinfo(np.int16).min and c_max < np.iinfo(np.int16).max:\n                    df[col] = df[col].astype(np.int16)\n                elif c_min > np.iinfo(np.int32).min and c_max < np.iinfo(np.int32).max:\n                    df[col] = df[col].astype(np.int32)\n            else:\n                if c_min > np.finfo(np.float32).min and c_max < np.finfo(np.float32).max:\n                    df[col] = df[col].astype(np.float32)\n    \n    end_mem = df.memory_usage().sum() / 1024**2\n    print(f'Memory usage after optimization: {end_mem:.2f} MB ({100 * (start_mem - end_mem) / start_mem:.1f}% reduction)')\n    gc.collect()\n    return df\n\ndef create_base_features(df):\n    \"\"\"Create market microstructure features with additional complexity\"\"\"\n    features_created = []\n    \n    if 'bid_qty' in df.columns and 'ask_qty' in df.columns:\n        df['bid_ask_spread'] = df['ask_qty'] - df['bid_qty']\n        df['bid_ask_imbalance'] = safe_divide(\n            df['bid_qty'] - df['ask_qty'],\n            df['bid_qty'] + df['ask_qty']\n        )\n        df['total_depth'] = df['bid_qty'] + df['ask_qty']\n        df['bid_ratio'] = safe_divide(df['bid_qty'], df['total_depth'])\n        df['depth_imbalance_squared'] = df['bid_ask_imbalance'] ** 2\n        features_created.extend(['bid_ask_spread', 'bid_ask_imbalance', 'total_depth', 'bid_ratio', 'depth_imbalance_squared'])\n    \n    if 'buy_qty' in df.columns and 'sell_qty' in df.columns:\n        df['order_imbalance'] = safe_divide(\n            df['buy_qty'] - df['sell_qty'],\n            df['buy_qty'] + df['sell_qty']\n        )\n        df['order_flow'] = df['buy_qty'] - df['sell_qty']\n        df['trade_intensity'] = df['buy_qty'] + df['sell_qty']\n        df['buy_ratio'] = safe_divide(df['buy_qty'], df['trade_intensity'])\n        df['order_flow_squared'] = df['order_flow'] ** 2\n        df['trade_intensity_log'] = np.log1p(df['trade_intensity'])\n        features_created.extend(['order_imbalance', 'order_flow', 'trade_intensity', 'buy_ratio', 'order_flow_squared', 'trade_intensity_log'])\n    \n    if 'volume' in df.columns:\n        df['log_volume'] = np.log1p(df['volume'])\n        df['sqrt_volume'] = np.sqrt(df['volume'])\n        df['volume_squared'] = df['volume'] ** 2\n        features_created.extend(['log_volume', 'sqrt_volume', 'volume_squared'])\n    \n    # Cross features between market data\n    if all(col in df.columns for col in ['bid_qty', 'buy_qty']):\n        df['bid_buy_ratio'] = safe_divide(df['bid_qty'], df['buy_qty'])\n        df['bid_buy_interaction'] = df['bid_qty'] * df['buy_qty']\n        features_created.extend(['bid_buy_ratio', 'bid_buy_interaction'])\n    \n    if all(col in df.columns for col in ['ask_qty', 'sell_qty']):\n        df['ask_sell_ratio'] = safe_divide(df['ask_qty'], df['sell_qty'])\n        df['ask_sell_interaction'] = df['ask_qty'] * df['sell_qty']\n        features_created.extend(['ask_sell_ratio', 'ask_sell_interaction'])\n    \n    # Complex microstructure features\n    if all(col in df.columns for col in ['bid_qty', 'ask_qty', 'volume']):\n        df['volume_to_depth_ratio'] = safe_divide(df['volume'], df['total_depth'])\n        df['spread_volume_interaction'] = df['bid_ask_spread'] * df['volume']\n        features_created.extend(['volume_to_depth_ratio', 'spread_volume_interaction'])\n    \n    return features_created\n\ndef clean_features(df, features):\n    \"\"\"Clean feature data with enhanced handling\"\"\"\n    df_clean = df[features].copy()\n    \n    # Replace infinities\n    df_clean = df_clean.replace([np.inf, -np.inf], np.nan)\n    \n    # Fill missing values with median\n    for col in features:\n        if col in df_clean.columns:\n            median_val = df_clean[col].median()\n            if pd.isna(median_val):\n                median_val = 0\n            df_clean[col] = df_clean[col].fillna(median_val)\n            \n            # Clip extreme values\n            percentile_99 = df_clean[col].quantile(0.99)\n            percentile_1 = df_clean[col].quantile(0.01)\n            df_clean[col] = df_clean[col].clip(lower=percentile_1, upper=percentile_99)\n    \n    return df_clean\n\ndef calculate_feature_importance(train_data, features):\n    \"\"\"Calculate feature importance using enhanced LightGBM\"\"\"\n    print(\"\\nCalculating feature importance...\")\n    \n    # Sample data for speed\n    sample_size = min(100000, len(train_data))\n    sample_data = train_data.sample(n=sample_size, random_state=CFG.random_seed)\n    \n    X = clean_features(sample_data, features)\n    y = sample_data['label']\n    \n    # Enhanced LightGBM parameters\n    lgb_params = {\n        'objective': 'regression',\n        'metric': 'rmse',\n        'num_leaves': 63,\n        'learning_rate': 0.05,\n        'feature_fraction': 0.8,\n        'bagging_fraction': 0.8,\n        'bagging_freq': 5,\n        'verbose': -1,\n        'seed': CFG.random_seed,\n        'n_estimators': 200,\n        'reg_alpha': 0.1,\n        'reg_lambda': 0.1\n    }\n    \n    model = lgb.LGBMRegressor(**lgb_params)\n    model.fit(X, y)\n    \n    # Get importance\n    importance_df = pd.DataFrame({\n        'feature': features,\n        'importance': model.feature_importances_\n    }).sort_values('importance', ascending=False)\n    \n    print(f\"Top 15 important features:\")\n    for idx, row in importance_df.head(15).iterrows():\n        print(f\"  {row['feature']}: {row['importance']:.0f}\")\n    \n    return importance_df\n\ndef apply_transformation(series, transformation):\n    \"\"\"Apply mathematical transformation with extended options\"\"\"\n    if transformation == 'log':\n        return np.log1p(series - series.min() + 1)\n    elif transformation == 'sqrt':\n        return np.sqrt(series - series.min())\n    elif transformation == 'square':\n        return series ** 2\n    elif transformation == 'reciprocal':\n        return safe_divide(1, series)\n    elif transformation == 'exp':\n        clipped = np.clip(series, -10, 10)\n        return np.exp(clipped)\n    elif transformation == 'abs':\n        return np.abs(series)\n    elif transformation == 'tanh':\n        return np.tanh(series)\n    elif transformation == 'sigmoid':\n        return 1 / (1 + np.exp(-np.clip(series, -10, 10)))\n    return series\n\ndef evaluate_features_parallel(train_data, val_data, features_list):\n    \"\"\"Evaluate multiple feature sets in parallel manner\"\"\"\n    results = []\n    \n    for features in features_list:\n        try:\n            score = evaluate_features(train_data, val_data, features)\n            results.append(score)\n        except Exception as e:\n            print(f\"Error evaluating features: {e}\")\n            results.append(-1)\n    \n    return results\n\ndef evaluate_features(train_data, val_data, features):\n    \"\"\"Evaluate feature set performance with enhanced model\"\"\"\n    # Prepare data\n    X_train = clean_features(train_data, features)\n    y_train = train_data['label']\n    X_val = clean_features(val_data, features)\n    y_val = val_data['label']\n    \n    # Enhanced model parameters\n    lgb_params = {\n        'objective': 'regression',\n        'metric': 'rmse',\n        'num_leaves': 63,\n        'learning_rate': 0.05,\n        'verbose': -1,\n        'seed': CFG.random_seed,\n        'n_estimators': 100,\n        'reg_alpha': 0.1,\n        'reg_lambda': 0.1\n    }\n    \n    model = lgb.LGBMRegressor(**lgb_params)\n    model.fit(X_train, y_train, eval_set=[(X_val, y_val)], callbacks=[lgb.early_stopping(20)])\n    \n    # Evaluate\n    val_pred = model.predict(X_val)\n    score = pearsonr(y_val, val_pred)[0]\n    \n    return score\n\ndef generate_feature_candidates(features, importance_dict, width, existing_features, depth):\n    \"\"\"Generate feature candidates with progressive complexity\"\"\"\n    candidates = []\n    \n    # Sort by importance\n    sorted_features = sorted(features, key=lambda x: importance_dict.get(x, 0), reverse=True)\n    \n    # Adjust selection based on depth\n    top_n = min(50 + depth * 5, len(sorted_features))\n    top_features = sorted_features[:top_n]\n    \n    # Transformations\n    n_transformations = int(width * CFG.transformation_weight)\n    transform_features = top_features[:min(25 + depth * 2, len(top_features))]\n    \n    for i, feat in enumerate(transform_features):\n        for transform in CFG.transformations:\n            if len(candidates) >= n_transformations:\n                break\n            feature_name = f\"{feat}_{transform}\"\n            if feature_name not in existing_features:\n                candidates.append({\n                    'type': 'transformation',\n                    'feature': feat,\n                    'transformation': transform\n                })\n    \n    # Binary combinations\n    operations = ['multiply', 'divide', 'add', 'subtract']\n    combination_features = top_features[:min(30 + depth * 2, len(top_features))]\n    \n    for feat1, feat2 in itertools.combinations(combination_features, 2):\n        if len(candidates) >= width * 0.7:\n            break\n        for op in operations:\n            if len(candidates) >= width * 0.7:\n                break\n            feature_name = f\"{feat1}_{op[:3]}_{feat2}\"\n            if feature_name not in existing_features:\n                candidates.append({\n                    'type': 'combination',\n                    'feat1': feat1,\n                    'feat2': feat2,\n                    'operation': op\n                })\n    \n    # Advanced combinations\n    if len(candidates) < width:\n        for op in CFG.advanced_operations:\n            for feat1, feat2 in itertools.combinations(top_features[:20], 2):\n                if len(candidates) >= width:\n                    break\n                feature_name = f\"{feat1}_{op}_{feat2}\"\n                if feature_name not in existing_features:\n                    candidates.append({\n                        'type': 'advanced',\n                        'feat1': feat1,\n                        'feat2': feat2,\n                        'operation': op\n                    })\n    \n    # Three-way interactions for later depths\n    if depth >= 5 and len(candidates) < width:\n        for feat1, feat2, feat3 in itertools.combinations(top_features[:10], 3):\n            if len(candidates) >= width:\n                break\n            feature_name = f\"{feat1}_x_{feat2}_x_{feat3}\"\n            if feature_name not in existing_features:\n                candidates.append({\n                    'type': 'triple',\n                    'feat1': feat1,\n                    'feat2': feat2,\n                    'feat3': feat3\n                })\n    \n    return candidates[:width]\n\ndef create_new_feature(df, candidate):\n    \"\"\"Create feature from candidate with extended operations\"\"\"\n    if candidate['type'] == 'transformation':\n        feature_name = f\"{candidate['feature']}_{candidate['transformation']}\"\n        values = apply_transformation(df[candidate['feature']], candidate['transformation'])\n        \n    elif candidate['type'] == 'combination':\n        feat1 = candidate['feat1']\n        feat2 = candidate['feat2']\n        op = candidate['operation']\n        \n        if op == 'multiply':\n            feature_name = f\"{feat1}_mul_{feat2}\"\n            values = df[feat1] * df[feat2]\n        elif op == 'divide':\n            feature_name = f\"{feat1}_div_{feat2}\"\n            values = safe_divide(df[feat1], df[feat2])\n        elif op == 'add':\n            feature_name = f\"{feat1}_add_{feat2}\"\n            values = df[feat1] + df[feat2]\n        else:  # subtract\n            feature_name = f\"{feat1}_sub_{feat2}\"\n            values = df[feat1] - df[feat2]\n    \n    elif candidate['type'] == 'advanced':\n        feat1 = candidate['feat1']\n        feat2 = candidate['feat2']\n        op = candidate['operation']\n        \n        if op == 'max':\n            feature_name = f\"{feat1}_max_{feat2}\"\n            values = np.maximum(df[feat1], df[feat2])\n        elif op == 'min':\n            feature_name = f\"{feat1}_min_{feat2}\"\n            values = np.minimum(df[feat1], df[feat2])\n        elif op == 'mean':\n            feature_name = f\"{feat1}_mean_{feat2}\"\n            values = (df[feat1] + df[feat2]) / 2\n        elif op == 'std':\n            feature_name = f\"{feat1}_std_{feat2}\"\n            values = np.sqrt(((df[feat1] - df[feat1].mean())**2 + (df[feat2] - df[feat2].mean())**2) / 2)\n        else:  # weighted_mean\n            feature_name = f\"{feat1}_wmean_{feat2}\"\n            weight1 = np.abs(df[feat1]).sum()\n            weight2 = np.abs(df[feat2]).sum()\n            total_weight = weight1 + weight2\n            if total_weight > 0:\n                values = (df[feat1] * weight1 + df[feat2] * weight2) / total_weight\n            else:\n                values = (df[feat1] + df[feat2]) / 2\n    \n    elif candidate['type'] == 'triple':\n        feat1 = candidate['feat1']\n        feat2 = candidate['feat2']\n        feat3 = candidate['feat3']\n        feature_name = f\"{feat1}_x_{feat2}_x_{feat3}\"\n        values = df[feat1] * df[feat2] * df[feat3]\n    \n    return feature_name, values\n\ndef greedy_feature_engineering(train_data, initial_features, importance_df, pipeline_state):\n    \"\"\"Perform greedy feature engineering with comprehensive checkpointing\"\"\"\n    print(\"\\nStarting extended feature engineering...\")\n    \n    # Check for existing checkpoint\n    checkpoint_state = load_checkpoint(CFG.checkpoint_features)\n    checkpoint_history = load_checkpoint(CFG.checkpoint_history)\n    \n    if checkpoint_state and checkpoint_history:\n        print(f\"Resuming from checkpoint with {len(checkpoint_state['features'])} features\")\n        current_features = checkpoint_state['features']\n        best_score = checkpoint_state['best_score']\n        feature_history = checkpoint_history\n        start_depth = checkpoint_state.get('last_depth', 0) + 1\n        last_checkpoint_time = datetime.fromisoformat(checkpoint_state.get('last_checkpoint_time', datetime.now().isoformat()))\n        \n        # Recreate engineered features\n        print(\"Recreating engineered features...\")\n        for i, feat_info in enumerate(feature_history):\n            if i % 10 == 0:\n                print(f\"  Progress: {i}/{len(feature_history)} features recreated\")\n            try:\n                candidate = feat_info['details']\n                feature_name, train_values = create_new_feature(train_data, candidate)\n                if feature_name not in train_data.columns:\n                    train_data[feature_name] = train_values\n            except Exception as e:\n                print(f\"Error recreating feature {feat_info.get('feature', 'unknown')}: {e}\")\n    else:\n        current_features = initial_features.copy()\n        feature_history = []\n        start_depth = 0\n        best_score = None\n        last_checkpoint_time = datetime.now()\n    \n    # Create validation split\n    val_size = int(len(train_data) * 0.2)\n    val_indices = np.random.choice(len(train_data), size=val_size, replace=False)\n    val_data = train_data.iloc[val_indices].copy()\n    train_fe = train_data.drop(train_data.index[val_indices]).copy()\n    \n    # Calculate baseline if needed\n    if best_score is None:\n        best_score = evaluate_features(train_fe, val_data, current_features)\n        print(f\"Baseline score: {best_score:.6f}\")\n    \n    # Get importance dict\n    importance_dict = dict(zip(importance_df['feature'], importance_df['importance']))\n    \n    # Track all created features to avoid duplicates\n    all_created_features = set(current_features)\n    \n    # Feature engineering loop\n    features_added_since_checkpoint = 0\n    \n    for depth in range(start_depth, CFG.new_feature_depth):\n        # Check runtime limit\n        exceeded, remaining_hours = check_runtime_limit(pipeline_state.start_time)\n        if exceeded:\n            print(\"Runtime limit exceeded, saving progress and exiting...\")\n            break\n        \n        print(f\"\\nDepth {depth + 1}/{CFG.new_feature_depth}\")\n        print(f\"Current features: {len(current_features)}\")\n        print(f\"Memory usage: {get_memory_usage():.2f} GB\")\n        \n        # Update pipeline state\n        pipeline_state.substage = f\"feature_engineering_depth_{depth + 1}\"\n        pipeline_state.progress['current_depth'] = depth + 1\n        pipeline_state.progress['features_count'] = len(current_features)\n        \n        # Generate candidates with depth-aware strategy\n        candidates = generate_feature_candidates(\n            current_features, importance_dict, CFG.feature_width, all_created_features, depth\n        )\n        print(f\"Evaluating {len(candidates)} candidates...\")\n        \n        best_candidate = None\n        best_new_score = best_score\n        improvement_found = False\n        candidates_evaluated = 0\n        \n        # Batch evaluation for efficiency\n        batch_size = 10\n        \n        for batch_start in range(0, len(candidates), batch_size):\n            batch_end = min(batch_start + batch_size, len(candidates))\n            batch_candidates = candidates[batch_start:batch_end]\n            \n            if batch_start % 20 == 0:\n                print(f\"  Progress: {batch_start}/{len(candidates)} candidates evaluated\")\n            \n            for candidate in batch_candidates:\n                try:\n                    # Create feature\n                    feature_name, train_values = create_new_feature(train_fe, candidate)\n                    \n                    if feature_name in train_fe.columns:\n                        continue\n                    \n                    # Add temporarily\n                    train_fe[feature_name] = train_values\n                    val_data[feature_name] = create_new_feature(val_data, candidate)[1]\n                    \n                    # Evaluate\n                    temp_features = current_features + [feature_name]\n                    score = evaluate_features(train_fe, val_data, temp_features)\n                    \n                    if score > best_new_score:\n                        best_candidate = candidate\n                        best_candidate['name'] = feature_name\n                        best_new_score = score\n                        improvement_found = True\n                    \n                    # Remove\n                    train_fe.drop(columns=[feature_name], inplace=True)\n                    val_data.drop(columns=[feature_name], inplace=True)\n                    \n                    candidates_evaluated += 1\n                    \n                except Exception as e:\n                    print(f\"  Error evaluating candidate: {e}\")\n                    continue\n            \n            # Periodic memory optimization\n            if candidates_evaluated % CFG.memory_optimization_interval == 0:\n                gc.collect()\n        \n        # Add best feature\n        if best_candidate and (best_new_score - best_score) > CFG.min_improvement:\n            feature_name = best_candidate['name']\n            train_fe[feature_name] = create_new_feature(train_fe, best_candidate)[1]\n            val_data[feature_name] = create_new_feature(val_data, best_candidate)[1]\n            \n            current_features.append(feature_name)\n            all_created_features.add(feature_name)\n            improvement = best_new_score - best_score\n            best_score = best_new_score\n            \n            print(f\"Added {feature_name}: {best_score:.6f} (+{improvement:.6f})\")\n            \n            feature_history.append({\n                'feature': feature_name,\n                'type': best_candidate['type'],\n                'details': best_candidate,\n                'score': best_score,\n                'improvement': improvement,\n                'depth': depth + 1,\n                'timestamp': datetime.now().isoformat()\n            })\n            \n            features_added_since_checkpoint += 1\n            \n            # Save checkpoint after interval or time threshold\n            time_since_checkpoint = (datetime.now() - last_checkpoint_time).total_seconds() / 60\n            \n            if (features_added_since_checkpoint >= CFG.feature_checkpoint_interval or \n                time_since_checkpoint >= CFG.checkpoint_interval_minutes):\n                \n                checkpoint_state = {\n                    'features': current_features,\n                    'best_score': best_score,\n                    'last_depth': depth,\n                    'last_checkpoint_time': datetime.now().isoformat(),\n                    'total_features_engineered': len(feature_history)\n                }\n                save_checkpoint(CFG.checkpoint_features, checkpoint_state)\n                save_checkpoint(CFG.checkpoint_history, feature_history)\n                save_pipeline_state(pipeline_state)\n                \n                features_added_since_checkpoint = 0\n                last_checkpoint_time = datetime.now()\n                print(f\"Checkpoint saved at depth {depth + 1}\")\n            \n        else:\n            print(\"No improvement found at this depth\")\n            if not improvement_found:\n                print(\"Early stopping - no candidates showed improvement\")\n                # Save final state before stopping\n                checkpoint_state = {\n                    'features': current_features,\n                    'best_score': best_score,\n                    'last_depth': depth,\n                    'last_checkpoint_time': datetime.now().isoformat(),\n                    'total_features_engineered': len(feature_history)\n                }\n                save_checkpoint(CFG.checkpoint_features, checkpoint_state)\n                save_checkpoint(CFG.checkpoint_history, feature_history)\n                break\n    \n    return current_features, feature_history\n\ndef train_final_model(train_data, test_data, features, pipeline_state):\n    \"\"\"Train final model with extended time budget and checkpointing\"\"\"\n    print(\"\\nTraining final model with extended parameters...\")\n    print(f\"Time budget: {CFG.final_model_time/3600:.1f} hours\")\n    \n    # Check for existing model checkpoint\n    model_checkpoint = load_checkpoint(CFG.checkpoint_model)\n    if model_checkpoint:\n        print(\"Found existing model checkpoint\")\n        test_predictions = model_checkpoint.get('predictions')\n        cv_scores = model_checkpoint.get('cv_scores', [])\n        if test_predictions is not None and len(cv_scores) > 0:\n            print(f\"Using saved predictions with CV score: {np.mean(cv_scores):.6f}\")\n            return test_predictions, cv_scores\n    \n    # Prepare data\n    X_train = clean_features(train_data, features)\n    y_train = train_data['label']\n    X_test = clean_features(test_data, features)\n    \n    # Update pipeline state\n    pipeline_state.stage = \"model_training\"\n    \n    # Train with FLAML\n    try:\n        automl = AutoML()\n        automl_settings = {\n            \"time_budget\": CFG.final_model_time,\n            \"metric\": 'r2',\n            \"estimator_list\": ['lgbm', 'xgboost', 'catboost', 'rf', 'extra_tree'],  # Added more estimators\n            \"task\": \"regression\",\n            \"seed\": CFG.random_seed,\n            \"verbose\": 2,  # Increased verbosity\n            \"eval_method\": \"cv\",\n            \"n_splits\": 5,\n            \"early_stop\": True,\n            \"retrain_full\": True,\n            \"ensemble\": True,  # Enable ensemble\n            \"max_iter\": 10000,  # Increased iterations\n            \"mem_thres\": 128 * 1024 * 1024 * 1024  # 128GB memory threshold\n        }\n        \n        print(\"Starting extended AutoML training...\")\n        automl.fit(X_train, y_train, **automl_settings)\n        test_predictions = automl.predict(X_test)\n        \n        print(f\"Best estimator: {automl.best_estimator}\")\n        print(f\"Best config: {automl.best_config}\")\n        print(f\"Best validation score: {automl.best_loss}\")\n        \n    except Exception as e:\n        # Enhanced fallback\n        print(f\"AutoML error: {e}\")\n        print(\"Using enhanced LightGBM ensemble fallback...\")\n        \n        # Train multiple models with different parameters\n        models = []\n        model_configs = [\n            {'num_leaves': 127, 'learning_rate': 0.03, 'n_estimators': 2000},\n            {'num_leaves': 63, 'learning_rate': 0.05, 'n_estimators': 1500},\n            {'num_leaves': 255, 'learning_rate': 0.02, 'n_estimators': 1000},\n        ]\n        \n        for config in model_configs:\n            lgb_params = {\n                'objective': 'regression',\n                'metric': 'rmse',\n                'feature_fraction': 0.8,\n                'bagging_fraction': 0.8,\n                'bagging_freq': 5,\n                'verbose': -1,\n                'seed': CFG.random_seed,\n                'reg_alpha': 0.1,\n                'reg_lambda': 0.1,\n                **config\n            }\n            \n            model = lgb.LGBMRegressor(**lgb_params)\n            model.fit(X_train, y_train)\n            models.append(model)\n        \n        # Ensemble predictions\n        predictions = []\n        for model in models:\n            predictions.append(model.predict(X_test))\n        test_predictions = np.mean(predictions, axis=0)\n    \n    # Enhanced cross-validation\n    kf = KFold(n_splits=CFG.n_folds, shuffle=True, random_state=CFG.random_seed)\n    cv_scores = []\n    \n    print(\"Calculating detailed cross-validation scores...\")\n    for fold, (train_idx, val_idx) in enumerate(kf.split(X_train)):\n        print(f\"  Processing fold {fold + 1}/{CFG.n_folds}\")\n        \n        X_fold_train = X_train.iloc[train_idx]\n        y_fold_train = y_train.iloc[train_idx]\n        X_fold_val = X_train.iloc[val_idx]\n        y_fold_val = y_train.iloc[val_idx]\n        \n        # Use best parameters from AutoML if available\n        if 'automl' in locals() and hasattr(automl, 'best_config'):\n            fold_model = lgb.LGBMRegressor(**automl.best_config)\n        else:\n            fold_model = lgb.LGBMRegressor(\n                n_estimators=500, \n                num_leaves=63,\n                learning_rate=0.05,\n                random_state=CFG.random_seed, \n                verbose=-1\n            )\n        \n        fold_model.fit(X_fold_train, y_fold_train)\n        fold_pred = fold_model.predict(X_fold_val)\n        \n        fold_score = pearsonr(y_fold_val, fold_pred)[0]\n        cv_scores.append(fold_score)\n        print(f\"    Fold {fold + 1} score: {fold_score:.6f}\")\n    \n    print(f\"CV Score: {np.mean(cv_scores):.6f} (±{np.std(cv_scores):.6f})\")\n    \n    # Save model checkpoint\n    model_checkpoint = {\n        'predictions': test_predictions,\n        'cv_scores': cv_scores,\n        'timestamp': datetime.now().isoformat(),\n        'features_used': len(features)\n    }\n    save_checkpoint(CFG.checkpoint_model, model_checkpoint)\n    \n    return test_predictions, cv_scores\n\ndef main():\n    \"\"\"Main execution with comprehensive checkpoint support\"\"\"\n    # Load or create pipeline state\n    pipeline_state = load_pipeline_state()\n    \n    if pipeline_state.start_time is None:\n        pipeline_state.start_time = datetime.now()\n        timestamp = pipeline_state.start_time.strftime(\"%Y%m%d_%H%M%S\")\n    else:\n        timestamp = pipeline_state.start_time.strftime(\"%Y%m%d_%H%M%S\")\n        print(f\"Resuming pipeline from {pipeline_state.stage or 'beginning'}\")\n    \n    print(\"=\"*80)\n    print(\"DRW CRYPTO MARKET PREDICTION - EXTENDED RUNTIME PIPELINE\")\n    print(f\"Timestamp: {timestamp}\")\n    print(f\"Maximum runtime: {CFG.max_runtime_hours} hours\")\n    print(\"=\"*80)\n    \n    # Configuration summary\n    print(\"\\nExtended Configuration:\")\n    print(f\"  Feature engineering depth: {CFG.new_feature_depth}\")\n    print(f\"  Candidates per iteration: {CFG.feature_width}\")\n    print(f\"  Minimum improvement threshold: {CFG.min_improvement}\")\n    print(f\"  Feature evaluation time: {CFG.feature_evaluation_time}s\")\n    print(f\"  Final model time budget: {CFG.final_model_time/3600:.1f} hours\")\n    print(f\"  Checkpoint interval: {CFG.checkpoint_interval_minutes} minutes\")\n    \n    # Stage 1: Data Loading\n    if pipeline_state.stage in [None, 'data_loading']:\n        pipeline_state.stage = 'data_loading'\n        print(\"\\nStage 1: Loading data...\")\n        \n        data_checkpoint = load_checkpoint(CFG.checkpoint_data)\n        if data_checkpoint:\n            print(\"Loading data from checkpoint...\")\n            train_data = data_checkpoint['train']\n            test_data = data_checkpoint['test']\n            sample_submission = data_checkpoint['submission']\n        else:\n            train_data = pd.read_parquet(CFG.train_path)\n            test_data = pd.read_parquet(CFG.test_path)\n            sample_submission = pd.read_csv(CFG.sample_sub_path)\n            \n            # Save data checkpoint\n            data_checkpoint = {\n                'train': train_data,\n                'test': test_data,\n                'submission': sample_submission\n            }\n            save_checkpoint(CFG.checkpoint_data, data_checkpoint)\n        \n        print(f\"Train shape: {train_data.shape}\")\n        print(f\"Test shape: {test_data.shape}\")\n        \n        # Optimize memory\n        train_data = reduce_memory_usage(train_data, \"train\")\n        test_data = reduce_memory_usage(test_data, \"test\")\n        \n        pipeline_state.results['data_loaded'] = True\n        save_pipeline_state(pipeline_state)\n    else:\n        # Load from checkpoint\n        data_checkpoint = load_checkpoint(CFG.checkpoint_data)\n        train_data = data_checkpoint['train']\n        test_data = data_checkpoint['test']\n        sample_submission = data_checkpoint['submission']\n    \n    # Stage 2: Feature Preparation\n    if pipeline_state.stage in [None, 'data_loading', 'feature_preparation']:\n        pipeline_state.stage = 'feature_preparation'\n        print(\"\\nStage 2: Feature preparation...\")\n        \n        # Get features\n        x_features = [col for col in train_data.columns if col.startswith('X')]\n        market_features = [\"bid_qty\", \"ask_qty\", \"buy_qty\", \"sell_qty\", \"volume\"]\n        \n        # Create base features\n        base_features = create_base_features(train_data)\n        test_base = create_base_features(test_data)\n        \n        # All features\n        all_features = x_features + market_features + base_features\n        all_features = [f for f in all_features if f in train_data.columns and f in test_data.columns]\n        \n        print(f\"\\nTotal features: {len(all_features)}\")\n        \n        # Feature importance\n        importance_df = calculate_feature_importance(train_data, all_features)\n        \n        # Select top features\n        initial_features = importance_df.head(CFG.max_base_features)['feature'].tolist()\n        \n        pipeline_state.results['initial_features'] = initial_features\n        pipeline_state.results['importance_df'] = importance_df\n        save_pipeline_state(pipeline_state)\n    else:\n        initial_features = pipeline_state.results['initial_features']\n        importance_df = pipeline_state.results['importance_df']\n    \n    # Stage 3: Feature Engineering\n    if pipeline_state.stage in [None, 'data_loading', 'feature_preparation', 'feature_engineering']:\n        pipeline_state.stage = 'feature_engineering'\n        \n        # Check runtime before starting expensive feature engineering\n        exceeded, remaining_hours = check_runtime_limit(pipeline_state.start_time)\n        if not exceeded and remaining_hours > 0.5:  # Need at least 30 minutes\n            final_features, feature_history = greedy_feature_engineering(\n                train_data, initial_features, importance_df, pipeline_state\n            )\n            \n            pipeline_state.results['final_features'] = final_features\n            pipeline_state.results['feature_history'] = feature_history\n            save_pipeline_state(pipeline_state)\n        else:\n            print(\"Insufficient time for feature engineering, using existing features\")\n            final_features = pipeline_state.results.get('final_features', initial_features)\n            feature_history = pipeline_state.results.get('feature_history', [])\n    else:\n        final_features = pipeline_state.results['final_features']\n        feature_history = pipeline_state.results['feature_history']\n    \n    # Apply engineered features to full dataset\n    print(\"\\nApplying engineered features to full dataset...\")\n    features_applied = 0\n    for feat_info in feature_history:\n        try:\n            candidate = feat_info['details']\n            feature_name, train_values = create_new_feature(train_data, candidate)\n            _, test_values = create_new_feature(test_data, candidate)\n            \n            if feature_name not in train_data.columns:\n                train_data[feature_name] = train_values\n            if feature_name not in test_data.columns:\n                test_data[feature_name] = test_values\n            \n            features_applied += 1\n            if features_applied % 20 == 0:\n                print(f\"  Applied {features_applied}/{len(feature_history)} features\")\n                \n        except Exception as e:\n            print(f\"Error applying feature {feat_info['feature']}: {e}\")\n            final_features = [f for f in final_features if f != feat_info['feature']]\n    \n    # Ensure all features exist\n    final_features = [f for f in final_features if f in train_data.columns and f in test_data.columns]\n    \n    print(f\"\\nFinal features: {len(final_features)}\")\n    print(f\"Engineered features added: {len(feature_history)}\")\n    \n    # Stage 4: Model Training\n    exceeded, remaining_hours = check_runtime_limit(pipeline_state.start_time)\n    if not exceeded and remaining_hours > 0.25:  # Need at least 15 minutes\n        test_predictions, cv_scores = train_final_model(train_data, test_data, final_features, pipeline_state)\n    else:\n        print(\"Insufficient time for model training, using saved predictions if available\")\n        model_checkpoint = load_checkpoint(CFG.checkpoint_model)\n        if model_checkpoint:\n            test_predictions = model_checkpoint['predictions']\n            cv_scores = model_checkpoint.get('cv_scores', [])\n        else:\n            print(\"No saved predictions found, using simple model\")\n            # Quick fallback model\n            X_train = clean_features(train_data, final_features)\n            y_train = train_data['label']\n            X_test = clean_features(test_data, final_features)\n            \n            model = lgb.LGBMRegressor(n_estimators=100, random_state=CFG.random_seed)\n            model.fit(X_train, y_train)\n            test_predictions = model.predict(X_test)\n            cv_scores = [0.0]  # Placeholder\n    \n    # Save predictions\n    submission = sample_submission.copy()\n    submission['prediction'] = test_predictions\n    submission_filename = f\"submission_{timestamp}.csv\"\n    submission.to_csv(submission_filename, index=False)\n    \n    print(f\"\\nSaved predictions to {submission_filename}\")\n    \n    # Save comprehensive summary\n    runtime_hours = (datetime.now() - pipeline_state.start_time).total_seconds() / 3600\n    \n    summary = {\n        'timestamp': timestamp,\n        'runtime_hours': runtime_hours,\n        'configuration': {\n            'feature_depth': CFG.new_feature_depth,\n            'feature_width': CFG.feature_width,\n            'min_improvement': CFG.min_improvement,\n            'feature_evaluation_time': CFG.feature_evaluation_time,\n            'final_model_time': CFG.final_model_time,\n            'max_runtime_hours': CFG.max_runtime_hours\n        },\n        'features': {\n            'initial': len(initial_features),\n            'engineered': len(feature_history),\n            'final': len(final_features)\n        },\n        'cv_score': {\n            'mean': np.mean(cv_scores) if cv_scores else 0,\n            'std': np.std(cv_scores) if cv_scores else 0,\n            'scores': cv_scores\n        },\n        'feature_history': feature_history,\n        'top_engineered_features': [f['feature'] for f in sorted(\n            feature_history, key=lambda x: x.get('improvement', 0), reverse=True\n        )[:20]],\n        'pipeline_state': pipeline_state.to_dict()\n    }\n    \n    with open(f'summary_{timestamp}.json', 'w') as f:\n        json.dump(summary, f, indent=2, default=str)\n    \n    print(\"\\nPipeline completed successfully!\")\n    print(f\"Total runtime: {runtime_hours:.2f} hours\")\n    print(f\"Final CV Score: {np.mean(cv_scores):.6f}\")\n\nif __name__ == \"__main__\":\n    main()","metadata":{"trusted":true,"execution":{"iopub.status.busy":"2025-06-01T04:25:41.668421Z","iopub.execute_input":"2025-06-01T04:25:41.668704Z","iopub.status.idle":"2025-06-01T04:25:41.854101Z","shell.execute_reply.started":"2025-06-01T04:25:41.668688Z","shell.execute_reply":"2025-06-01T04:25:41.853098Z"}},"outputs":[],"execution_count":null}]}