{"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":"gpu","dataSources":[{"sourceId":96164,"databundleVersionId":11418275,"sourceType":"competition"}],"dockerImageVersionId":31041,"isInternetEnabled":true,"language":"python","sourceType":"notebook","isGpuEnabled":true}},"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,"jupyter":{"source_hidden":true},"execution":{"iopub.status.busy":"2025-06-04T16:05:27.465913Z","iopub.execute_input":"2025-06-04T16:05:27.466100Z","iopub.status.idle":"2025-06-04T16:05:29.171802Z","shell.execute_reply.started":"2025-06-04T16:05:27.466084Z","shell.execute_reply":"2025-06-04T16:05:29.171015Z"}},"outputs":[],"execution_count":null},{"cell_type":"code","source":"# INITIAL SETUP AND PACKAGE INSTALLATION\n# !/usr/bin/env python\n# -*- coding: utf-8 -*-\n\"\"\"\nDRW Crypto Market Prediction Pipeline - Setup and Global Configuration\nThis module provides the foundation for all prediction models with dynamic model discovery\n\"\"\"\n\nimport subprocess\nimport sys\nimport os\nimport gc\nimport warnings\nimport json\nfrom pathlib import Path\nfrom dataclasses import dataclass, field\nfrom typing import List, Dict, Tuple, Optional, Any\nfrom datetime import datetime\n\nwarnings.filterwarnings('ignore')\n\n# Install only essential base packages\nprint(\"Installing base packages...\")\nbase_packages = [\"pandas\", \"numpy\", \"scipy\", \"scikit-learn\"]\nsubprocess.check_call([sys.executable, \"-m\", \"pip\", \"install\"] + base_packages + [\"--quiet\"])\n\n# Import base packages after installation\nimport pandas as pd\nimport numpy as np\nfrom scipy.stats import pearsonr\nfrom sklearn.model_selection import train_test_split, KFold, TimeSeriesSplit\nfrom sklearn.linear_model import Ridge\nfrom sklearn.base import clone, BaseEstimator, RegressorMixin\nfrom sklearn.metrics import mean_absolute_error\n\n# Memory management function\ndef aggressive_memory_cleanup():\n    \"\"\"Aggressively clean up memory to prevent kernel death\"\"\"\n    # PyTorch cleanup\n    if 'torch' in sys.modules:\n        try:\n            import torch\n            if torch.cuda.is_available():\n                torch.cuda.empty_cache()\n                torch.cuda.synchronize()\n        except:\n            pass\n    \n    # TensorFlow cleanup\n    if 'tensorflow' in sys.modules:\n        try:\n            import tensorflow as tf\n            tf.keras.backend.clear_session()\n            if hasattr(tf.compat, 'v1'):\n                tf.compat.v1.reset_default_graph()\n        except:\n            pass\n    \n    # Remove loaded modules to free memory\n    modules_to_remove = ['xgboost', 'lightgbm', 'tensorflow', 'torch', 'sklearn', 'gplearn', 'catboost']\n    for module in list(sys.modules.keys()):\n        if any(module.startswith(mod) for mod in modules_to_remove):\n            try:\n                del sys.modules[module]\n            except:\n                pass\n    \n    # Multiple garbage collection passes\n    for _ in range(3):\n        gc.collect()\n\n# Configure compute environment\ndef configure_compute_environment():\n    \"\"\"Configure compute environment for optimal performance\"\"\"\n    os.environ['CUDA_VISIBLE_DEVICES'] = '-1'  # Default to CPU\n    os.environ['TF_CPP_MIN_LOG_LEVEL'] = '3'\n    os.environ['TF_FORCE_GPU_ALLOW_GROWTH'] = 'true'\n    os.environ['OMP_NUM_THREADS'] = '4'\n    os.environ['MKL_NUM_THREADS'] = '4'\n\n@dataclass\nclass ModelOutput:\n    \"\"\"Represents an output file from a model\"\"\"\n    file_path: str\n    file_type: str  # 'submission', 'analysis', 'features', etc.\n    created_at: datetime\n    metadata: Dict[str, Any] = field(default_factory=dict)\n\n@dataclass\nclass ModelRecord:\n    \"\"\"Complete record of a model's execution and outputs\"\"\"\n    name: str\n    directory: str\n    status: str = \"not_started\"\n    score: Optional[float] = None\n    error_message: Optional[str] = None\n    start_time: Optional[datetime] = None\n    end_time: Optional[datetime] = None\n    outputs: Dict[str, ModelOutput] = field(default_factory=dict)  # key: output_name, value: ModelOutput\n    config: Dict[str, Any] = field(default_factory=dict)\n    \n    def add_output(self, output_name: str, file_path: str, file_type: str = \"submission\", metadata: Optional[Dict] = None):\n        \"\"\"Register an output file from this model\"\"\"\n        self.outputs[output_name] = ModelOutput(\n            file_path=file_path,\n            file_type=file_type,\n            created_at=datetime.now(),\n            metadata=metadata or {}\n        )\n    \n    def get_submission_files(self) -> List[str]:\n        \"\"\"Get all submission files from this model\"\"\"\n        return [output.file_path for output in self.outputs.values() \n                if output.file_type == \"submission\" and os.path.exists(output.file_path)]\n    \n    def to_dict(self) -> Dict[str, Any]:\n        \"\"\"Convert to dictionary for serialization\"\"\"\n        return {\n            'name': self.name,\n            'directory': self.directory,\n            'status': self.status,\n            'score': self.score,\n            'error_message': self.error_message,\n            'start_time': self.start_time.isoformat() if self.start_time else None,\n            'end_time': self.end_time.isoformat() if self.end_time else None,\n            'outputs': {\n                name: {\n                    'file_path': output.file_path,\n                    'file_type': output.file_type,\n                    'created_at': output.created_at.isoformat(),\n                    'metadata': output.metadata\n                }\n                for name, output in self.outputs.items()\n            },\n            'config': self.config\n        }\n    \n    @classmethod\n    def from_dict(cls, data: Dict[str, Any]) -> 'ModelRecord':\n        \"\"\"Create from dictionary\"\"\"\n        # Convert time fields\n        if data.get('start_time'):\n            data['start_time'] = datetime.fromisoformat(data['start_time'])\n        if data.get('end_time'):\n            data['end_time'] = datetime.fromisoformat(data['end_time'])\n        \n        # Convert outputs\n        outputs = {}\n        for name, output_data in data.get('outputs', {}).items():\n            outputs[name] = ModelOutput(\n                file_path=output_data['file_path'],\n                file_type=output_data['file_type'],\n                created_at=datetime.fromisoformat(output_data['created_at']),\n                metadata=output_data.get('metadata', {})\n            )\n        data['outputs'] = outputs\n        \n        return cls(**data)\n\n@dataclass\nclass GlobalConfig:\n    \"\"\"Central configuration with dynamic model registry\"\"\"\n    # Base paths\n    base_dir: str = \"/kaggle/working/sub-models\"\n    train_path: str = \"/kaggle/input/drw-crypto-market-prediction/train.parquet\"\n    test_path: str = \"/kaggle/input/drw-crypto-market-prediction/test.parquet\"\n    sample_sub_path: str = \"/kaggle/input/drw-crypto-market-prediction/sample_submission.csv\"\n    \n    # Model registry - using ModelRecord for complete tracking\n    model_registry: Dict[str, ModelRecord] = field(default_factory=dict)\n    \n    # Pipeline metadata\n    pipeline_start_time: Optional[datetime] = None\n    pipeline_end_time: Optional[datetime] = None\n    \n    def __post_init__(self):\n        \"\"\"Initialize configuration\"\"\"\n        Path(self.base_dir).mkdir(parents=True, exist_ok=True)\n        self.load_state()\n    \n    def register_model(self, name: str, directory: Optional[str] = None, config: Optional[Dict] = None) -> ModelRecord:\n        \"\"\"Register a new model in the pipeline\"\"\"\n        if directory is None:\n            directory = os.path.join(self.base_dir, name)\n        \n        # Create directory if it doesn't exist\n        Path(directory).mkdir(parents=True, exist_ok=True)\n        \n        # Create or update model record\n        if name in self.model_registry:\n            model_record = self.model_registry[name]\n            model_record.directory = directory\n            if config:\n                model_record.config.update(config)\n        else:\n            model_record = ModelRecord(\n                name=name,\n                directory=directory,\n                config=config or {}\n            )\n            self.model_registry[name] = model_record\n        \n        self.save_state()\n        return model_record\n    \n    def register_model_output(self, model_name: str, output_name: str, file_path: str, \n                            file_type: str = \"submission\", metadata: Optional[Dict] = None):\n        \"\"\"Register an output file from a model\"\"\"\n        if model_name not in self.model_registry:\n            self.register_model(model_name)\n        \n        model_record = self.model_registry[model_name]\n        model_record.add_output(output_name, file_path, file_type, metadata)\n        self.save_state()\n    \n    def update_model_status(self, name: str, status: str, \n                           score: Optional[float] = None,\n                           error_message: Optional[str] = None):\n        \"\"\"Update the status of a model\"\"\"\n        if name not in self.model_registry:\n            self.register_model(name)\n        \n        model = self.model_registry[name]\n        model.status = status\n        \n        if status == \"running\":\n            model.start_time = datetime.now()\n        elif status in [\"completed\", \"failed\"]:\n            model.end_time = datetime.now()\n        \n        if score is not None:\n            model.score = score\n        \n        if error_message is not None:\n            model.error_message = error_message\n        \n        self.save_state()\n    \n    def get_model_submissions(self, model_name: str) -> List[str]:\n        \"\"\"Get all submission files for a specific model\"\"\"\n        if model_name not in self.model_registry:\n            return []\n        \n        return self.model_registry[model_name].get_submission_files()\n    \n    def get_all_submissions(self) -> Dict[str, List[str]]:\n        \"\"\"Get all available submission files from all models\"\"\"\n        submissions = {}\n        for name, model in self.model_registry.items():\n            if model.status == \"completed\":\n                submission_files = model.get_submission_files()\n                if submission_files:\n                    submissions[name] = submission_files\n        return submissions\n    \n    def get_latest_submission(self, model_name: str) -> Optional[str]:\n        \"\"\"Get the most recent submission file from a model\"\"\"\n        submissions = self.get_model_submissions(model_name)\n        if not submissions:\n            return None\n        \n        # Get the output with the latest creation time\n        model = self.model_registry[model_name]\n        submission_outputs = [(name, output) for name, output in model.outputs.items() \n                            if output.file_type == \"submission\" and os.path.exists(output.file_path)]\n        \n        if not submission_outputs:\n            return None\n        \n        # Sort by creation time and return the latest\n        submission_outputs.sort(key=lambda x: x[1].created_at, reverse=True)\n        return submission_outputs[0][1].file_path\n    \n    def get_model_summary(self) -> pd.DataFrame:\n        \"\"\"Get a summary of all models as a DataFrame\"\"\"\n        data = []\n        for name, model in self.model_registry.items():\n            runtime = None\n            if model.start_time and model.end_time:\n                runtime = (model.end_time - model.start_time).total_seconds()\n            \n            submission_count = len(model.get_submission_files())\n            \n            data.append({\n                'Model': name,\n                'Status': model.status,\n                'Score': model.score,\n                'Runtime (s)': runtime,\n                'Directory': model.directory,\n                'Submissions': submission_count,\n                'Total Outputs': len(model.outputs),\n                'Error': model.error_message[:50] if model.error_message else None\n            })\n        \n        return pd.DataFrame(data)\n    \n    def get_execution_summary(self) -> str:\n        \"\"\"Get a formatted summary of model execution status\"\"\"\n        summary_df = self.get_model_summary()\n        summary_lines = []\n        \n        if summary_df.empty:\n            return \"No models registered yet.\"\n        \n        # Status counts\n        status_counts = summary_df['Status'].value_counts()\n        summary_lines.append(f\"Total models: {len(summary_df)}\")\n        \n        for status in ['completed', 'failed', 'running', 'not_started']:\n            count = status_counts.get(status, 0)\n            if count > 0:\n                emoji = {'completed': '✅', 'failed': '❌', 'running': '🔄', 'not_started': '⏸️'}[status]\n                summary_lines.append(f\"{emoji} {status}: {count}\")\n        \n        # Completed models with outputs\n        completed = summary_df[summary_df['Status'] == 'completed']\n        if not completed.empty:\n            summary_lines.append(\"\\nCompleted models:\")\n            for _, row in completed.iterrows():\n                score_str = f\"score: {row['Score']:.4f}\" if pd.notna(row['Score']) else \"score: N/A\"\n                outputs_str = f\"outputs: {row['Total Outputs']}\"\n                summary_lines.append(f\"  ✅ {row['Model']} ({score_str}, {outputs_str})\")\n        \n        # Failed models\n        failed = summary_df[summary_df['Status'] == 'failed']\n        if not failed.empty:\n            summary_lines.append(\"\\nFailed models:\")\n            for _, row in failed.iterrows():\n                error_str = row['Error'] if pd.notna(row['Error']) else \"Unknown error\"\n                summary_lines.append(f\"  ❌ {row['Model']}: {error_str}...\")\n        \n        return \"\\n\".join(summary_lines)\n    \n    def save_state(self):\n        \"\"\"Save current state to disk\"\"\"\n        state = {\n            'model_registry': {name: model.to_dict() for name, model in self.model_registry.items()},\n            'pipeline_start_time': self.pipeline_start_time.isoformat() if self.pipeline_start_time else None,\n            'pipeline_end_time': self.pipeline_end_time.isoformat() if self.pipeline_end_time else None\n        }\n        \n        state_path = os.path.join(self.base_dir, 'pipeline_state.json')\n        with open(state_path, 'w') as f:\n            json.dump(state, f, indent=2)\n        \n        # Save summary as CSV\n        summary_path = os.path.join(self.base_dir, 'model_summary.csv')\n        self.get_model_summary().to_csv(summary_path, index=False)\n    \n    def load_state(self):\n        \"\"\"Load state from disk if available\"\"\"\n        state_path = os.path.join(self.base_dir, 'pipeline_state.json')\n        \n        if os.path.exists(state_path):\n            try:\n                with open(state_path, 'r') as f:\n                    state = json.load(f)\n                \n                self.model_registry = {\n                    name: ModelRecord.from_dict(data) \n                    for name, data in state.get('model_registry', {}).items()\n                }\n                \n                if state.get('pipeline_start_time'):\n                    self.pipeline_start_time = datetime.fromisoformat(state['pipeline_start_time'])\n                if state.get('pipeline_end_time'):\n                    self.pipeline_end_time = datetime.fromisoformat(state['pipeline_end_time'])\n                    \n            except Exception as e:\n                print(f\"Warning: Could not load previous state: {e}\")\n    \n    def reset_pipeline(self):\n        \"\"\"Reset all models to initial state\"\"\"\n        self.model_registry.clear()\n        self.pipeline_start_time = None\n        self.pipeline_end_time = None\n        self.save_state()\n\n# Data utility functions\ndef reduce_memory_usage(df: pd.DataFrame, verbose: bool = True) -> pd.DataFrame:\n    \"\"\"Reduce memory usage by optimizing data types\"\"\"\n    if verbose:\n        start_mem = df.memory_usage().sum() / 1024**2\n        print(f\"Memory usage before: {start_mem:.2f} MB\")\n    \n    for col in df.columns:\n        col_type = df[col].dtype\n        \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.float16).min and c_max < np.finfo(np.float16).max:\n                    df[col] = df[col].astype(np.float16)\n                elif 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    if verbose:\n        end_mem = df.memory_usage().sum() / 1024**2\n        print(f\"Memory usage after: {end_mem:.2f} MB ({100 * (start_mem - end_mem) / start_mem:.1f}% reduction)\")\n    \n    return df\n\n# Initialize environment\nconfigure_compute_environment()\n\n# Initialize global configuration\nglobal_config = GlobalConfig()\n\nprint(\"🚀 DRW Crypto Prediction Pipeline - Setup Complete\")\nprint(\"=\"*80)\nprint(f\"Base directory: {global_config.base_dir}\")\nprint(f\"Training data: {global_config.train_path}\")\nprint(f\"Test data: {global_config.test_path}\")\nprint(f\"\\nModels will be automatically registered when they run.\")\nprint(\"Each model should register its outputs using global_config.register_model_output()\")\n\n# Display execution status if any\nif global_config.model_registry:\n    print(\"\\nPrevious execution status:\")\n    print(global_config.get_execution_summary())\nelse:\n    print(\"\\nNo previous execution history found - starting fresh\")\nprint(\"=\"*80)","metadata":{"trusted":true,"jupyter":{"outputs_hidden":true,"source_hidden":true},"execution":{"iopub.status.busy":"2025-06-04T16:05:29.173573Z","iopub.execute_input":"2025-06-04T16:05:29.173854Z","iopub.status.idle":"2025-06-04T16:05:34.002740Z","shell.execute_reply.started":"2025-06-04T16:05:29.173836Z","shell.execute_reply":"2025-06-04T16:05:34.001952Z"},"collapsed":true},"outputs":[],"execution_count":null},{"cell_type":"code","source":"# XGBOOST THREE-MODEL PIPELINE IMPLEMENTATION\n# !/usr/bin/env python\n# -*- coding: utf-8 -*-\n\"\"\"\nDRW Crypto Market Prediction - XGBoost Three-Model Pipeline\nThis module implements a time-weighted ensemble of three XGBoost models\ntrained on different time windows to capture both long-term and recent patterns\n\"\"\"\n\nimport subprocess\nimport sys\nimport pandas as pd\nimport numpy as np\nfrom scipy.stats import pearsonr\nfrom typing import List, Dict, Tuple, Optional, Any\nfrom pathlib import Path\nimport json\nimport gc\nimport warnings\nimport os\nfrom sklearn.model_selection import KFold\nimport matplotlib.pyplot as plt\n\nwarnings.filterwarnings('ignore')\n\n# Install required packages for this pipeline\nprint(\"Installing packages for XGBoost pipeline...\")\npackages_to_install = [\n    'xgboost==2.0.3',\n    'shap==0.44.0'\n]\n\nfor package in packages_to_install:\n    subprocess.check_call([sys.executable, \"-m\", \"pip\", \"install\", package, \"--quiet\"])\n\n# Import XGBoost after installation\nimport xgboost as xgb\n\nclass XGBoostConfiguration:\n    \"\"\"Configuration for XGBoost Three-Model Pipeline\"\"\"\n    \n    def __init__(self):\n        # Model registration\n        self.model_name = \"xgboost\"\n        self.model_directory = os.path.join(global_config.base_dir, \"triple_xgboost\")\n        \n        # Register model with global configuration\n        global_config.register_model(self.model_name, self.model_directory)\n        \n        # Data paths from global configuration\n        self.train_path = global_config.train_path\n        self.test_path = global_config.test_path\n        self.sample_sub_path = global_config.sample_sub_path\n        \n        # Model-specific feature selection\n        self.selected_features = [\n            \"X863\", \"X856\", \"X344\", \"X598\", \"X862\", \"X385\", \"X852\", \"X603\", \"X860\", \"X674\",\n            \"X415\", \"X345\", \"X137\", \"X855\", \"X174\", \"X302\", \"X178\", \"X532\", \"X168\", \"X612\",\n            \"bid_qty\", \"ask_qty\", \"buy_qty\", \"sell_qty\", \"volume\"\n        ]\n        \n        # XGBoost hyperparameters\n        self.xgb_params = {\n            \"tree_method\": \"hist\",  # Changed from gpu_hist for better compatibility\n            \"device\": \"cpu\",\n            \"colsample_bylevel\": 0.4778015829774066,\n            \"colsample_bynode\": 0.362764358742407,\n            \"colsample_bytree\": 0.7107423488010493,\n            \"gamma\": 1.7094857725240398,\n            \"learning_rate\": 0.02213323588455387,\n            \"max_depth\": 20,\n            \"max_leaves\": 12,\n            \"min_child_weight\": 16,\n            \"n_estimators\": 1667,\n            \"n_jobs\": -1,\n            \"random_state\": 42,\n            \"reg_alpha\": 39.352415706891264,\n            \"reg_lambda\": 75.44843704068275,\n            \"subsample\": 0.06566669853471274,\n            \"verbosity\": 0\n        }\n        \n        # Model configurations for time windows\n        self.model_configs = [\n            {\"name\": \"model_1_full_data\", \"percent\": 1.00, \"description\": \"Full Data\"},\n            {\"name\": \"model_2_recent_75\", \"percent\": 0.75, \"description\": \"75% Recent\"},\n            {\"name\": \"model_3_recent_50\", \"percent\": 0.50, \"description\": \"50% Recent\"}\n        ]\n        \n        # Cross-validation parameters\n        self.n_folds = 5\n        self.random_state = 42\n        self.shuffle = True\n        self.decay_factor = 0.95\n        self.early_stopping_rounds = 25\n        self.verbose_eval = 200\n        \n        # Output paths\n        self.intermediate_dir = os.path.join(self.model_directory, \"sub_models\")\n        self.submission_file = os.path.join(self.model_directory, \"submission.csv\")\n        self.results_file = os.path.join(self.model_directory, \"ensemble_results.csv\")\n        self.shap_features_path = os.path.join(self.model_directory, \"shap_features.csv\")\n        \n        # Ensure directories exist\n        Path(self.model_directory).mkdir(parents=True, exist_ok=True)\n        Path(self.intermediate_dir).mkdir(parents=True, exist_ok=True)\n\nclass XGBoostDataProcessor:\n    \"\"\"Data processing utilities for XGBoost pipeline\"\"\"\n    \n    @staticmethod\n    def create_time_weights(n_samples: int, decay_factor: float = 0.95) -> np.ndarray:\n        \"\"\"Create exponential decay weights for time series data\"\"\"\n        positions = np.arange(n_samples)\n        normalized_positions = positions / (n_samples - 1)\n        weights = decay_factor ** (1 - normalized_positions)\n        weights = weights * n_samples / weights.sum()\n        return weights\n    \n    @staticmethod\n    def load_data(config: XGBoostConfiguration) -> Tuple[pd.DataFrame, pd.DataFrame, pd.DataFrame]:\n        \"\"\"Load and prepare data for XGBoost models\"\"\"\n        print(\"Loading data...\")\n        train = pd.read_parquet(config.train_path).reset_index(drop=True)\n        test = pd.read_parquet(config.test_path).reset_index(drop=True)\n        sample = pd.read_csv(config.sample_sub_path)\n        \n        # Verify feature availability\n        available_features = [f for f in config.selected_features if f in train.columns]\n        missing_features = set(config.selected_features) - set(available_features)\n        \n        if missing_features:\n            print(f\"Warning: {len(missing_features)} features not found in data\")\n            print(f\"Missing features: {missing_features}\")\n        \n        # Select only available features\n        train = train[available_features + [\"label\"]]\n        test = test[available_features]\n        \n        # Reduce memory usage\n        train = reduce_memory_usage(train, verbose=False)\n        test = reduce_memory_usage(test, verbose=False)\n        \n        print(f\"Data loaded - Train: {train.shape}, Test: {test.shape}\")\n        print(f\"Using {len(available_features)} features\")\n        \n        return train, test, sample\n\nclass XGBoostModelTrainer:\n    \"\"\"Model training utilities for XGBoost pipeline\"\"\"\n    \n    def __init__(self, config: XGBoostConfiguration):\n        self.config = config\n    \n    def train_model(self, X_train: pd.DataFrame, y_train: pd.Series, \n                   X_valid: pd.DataFrame, y_valid: pd.Series, \n                   sample_weights: np.ndarray) -> xgb.XGBRegressor:\n        \"\"\"Train a single XGBoost model with early stopping\"\"\"\n        model = xgb.XGBRegressor(**self.config.xgb_params)\n        model.fit(\n            X_train, y_train,\n            sample_weight=sample_weights,\n            eval_set=[(X_valid, y_valid)],\n            early_stopping_rounds=self.config.early_stopping_rounds,\n            verbose=self.config.verbose_eval\n        )\n        return model\n    \n    def prepare_windowed_data(self, train_df: pd.DataFrame, train_idx: np.ndarray,\n                            cutoff: int, features: List[str]) -> Tuple[pd.DataFrame, pd.Series, np.ndarray]:\n        \"\"\"Prepare training data for windowed models\"\"\"\n        train_idx_recent = train_idx[train_idx >= cutoff]\n        train_idx_recent_adjusted = train_idx_recent - cutoff\n        \n        train_recent = train_df.iloc[cutoff:].reset_index(drop=True)\n        \n        X_train = train_recent.iloc[train_idx_recent_adjusted][features]\n        y_train = train_recent.iloc[train_idx_recent_adjusted][\"label\"]\n        \n        sample_weights_recent = XGBoostDataProcessor.create_time_weights(\n            len(train_recent), self.config.decay_factor\n        )\n        train_weights = sample_weights_recent[train_idx_recent_adjusted]\n        \n        return X_train, y_train, train_weights\n\nclass XGBoostEnsembleBuilder:\n    \"\"\"Ensemble building and evaluation utilities\"\"\"\n    \n    def __init__(self, config: XGBoostConfiguration):\n        self.config = config\n    \n    def evaluate_models(self, predictions: Dict[str, np.ndarray], \n                       train_labels: pd.Series) -> pd.DataFrame:\n        \"\"\"Evaluate individual models and ensemble combinations\"\"\"\n        scores = {}\n        \n        # Individual model scores\n        for model_name, preds in predictions.items():\n            scores[model_name] = pearsonr(train_labels, preds)[0]\n        \n        # Simple ensemble\n        simple_ensemble = np.mean(list(predictions.values()), axis=0)\n        scores['simple_ensemble'] = pearsonr(train_labels, simple_ensemble)[0]\n        \n        # Weighted ensemble based on individual scores\n        weights = np.array([scores[name] for name in predictions.keys()])\n        weights = weights / weights.sum()\n        \n        weighted_ensemble = np.average(list(predictions.values()), axis=0, weights=weights)\n        scores['weighted_ensemble'] = pearsonr(train_labels, weighted_ensemble)[0]\n        \n        # Create results dataframe\n        results_data = []\n        for i, (name, score) in enumerate(scores.items()):\n            weight = weights[i] if i < len(weights) else np.nan\n            results_data.append({\n                'model': name,\n                'pearson_correlation': score,\n                'weight_in_final': weight\n            })\n        \n        return pd.DataFrame(results_data)\n    \n    def create_final_ensemble(self, test_predictions: Dict[str, np.ndarray],\n                            oof_scores: Dict[str, float]) -> np.ndarray:\n        \"\"\"Create final ensemble predictions\"\"\"\n        # Use weighted ensemble if it performs better\n        weights = np.array([oof_scores[name] for name in test_predictions.keys()])\n        weights = weights / weights.sum()\n        \n        final_predictions = np.average(list(test_predictions.values()), axis=0, weights=weights)\n        \n        return final_predictions\n\nclass XGBoostThreeModelPipeline:\n    \"\"\"Main pipeline orchestrating the three-model XGBoost approach\"\"\"\n    \n    def __init__(self):\n        self.config = XGBoostConfiguration()\n        self.data_processor = XGBoostDataProcessor()\n        self.trainer = XGBoostModelTrainer(self.config)\n        self.ensemble_builder = XGBoostEnsembleBuilder(self.config)\n    \n    def run_cross_validation(self, train_df: pd.DataFrame, test_df: pd.DataFrame) -> Dict[str, Any]:\n        \"\"\"Execute cross-validation for all three models\"\"\"\n        features = [c for c in train_df.columns if c != \"label\"]\n        \n        # Calculate cutoffs for time windows\n        cutoff_75 = int(len(train_df) * 0.25)\n        cutoff_50 = int(len(train_df) * 0.50)\n        \n        print(f\"\\nTime window configurations:\")\n        print(f\"  Model 1: Full data ({len(train_df):,} samples)\")\n        print(f\"  Model 2: 75% recent ({len(train_df) - cutoff_75:,} samples)\")\n        print(f\"  Model 3: 50% recent ({len(train_df) - cutoff_50:,} samples)\")\n        \n        # Initialize prediction storage\n        oof_predictions = {\n            'model_1': np.zeros(len(train_df)),\n            'model_2': np.zeros(len(train_df)),\n            'model_3': np.zeros(len(train_df))\n        }\n        \n        test_predictions = {\n            'model_1': np.zeros(len(test_df)),\n            'model_2': np.zeros(len(test_df)),\n            'model_3': np.zeros(len(test_df))\n        }\n        \n        # Create sample weights for full dataset\n        sample_weights_full = self.data_processor.create_time_weights(\n            len(train_df), self.config.decay_factor\n        )\n        \n        # Cross-validation\n        kf = KFold(n_splits=self.config.n_folds, shuffle=self.config.shuffle, \n                  random_state=self.config.random_state)\n        \n        for fold, (train_idx, valid_idx) in enumerate(kf.split(train_df)):\n            print(f\"\\nProcessing fold {fold + 1}/{self.config.n_folds}\")\n            \n            X_valid = train_df.iloc[valid_idx][features]\n            y_valid = train_df.iloc[valid_idx][\"label\"]\n            X_test = test_df[features]\n            \n            # Model 1: Full data\n            X_train = train_df.iloc[train_idx][features]\n            y_train = train_df.iloc[train_idx][\"label\"]\n            train_weights = sample_weights_full[train_idx]\n            \n            model1 = self.trainer.train_model(X_train, y_train, X_valid, y_valid, train_weights)\n            oof_predictions['model_1'][valid_idx] = model1.predict(X_valid)\n            test_predictions['model_1'] += model1.predict(X_test) / self.config.n_folds\n            \n            # Model 2: 75% recent\n            X_train_75, y_train_75, weights_75 = self.trainer.prepare_windowed_data(\n                train_df, train_idx, cutoff_75, features\n            )\n            \n            if len(X_train_75) > 0:\n                model2 = self.trainer.train_model(X_train_75, y_train_75, X_valid, y_valid, weights_75)\n                \n                # Handle predictions for validation set\n                valid_idx_recent = valid_idx[valid_idx >= cutoff_75]\n                if len(valid_idx_recent) > 0:\n                    X_valid_recent = train_df.iloc[valid_idx_recent][features]\n                    oof_predictions['model_2'][valid_idx_recent] = model2.predict(X_valid_recent)\n                \n                # Use model 1 predictions for samples before cutoff\n                valid_idx_old = valid_idx[valid_idx < cutoff_75]\n                if len(valid_idx_old) > 0:\n                    oof_predictions['model_2'][valid_idx_old] = oof_predictions['model_1'][valid_idx_old]\n                \n                test_predictions['model_2'] += model2.predict(X_test) / self.config.n_folds\n            \n            # Model 3: 50% recent\n            X_train_50, y_train_50, weights_50 = self.trainer.prepare_windowed_data(\n                train_df, train_idx, cutoff_50, features\n            )\n            \n            if len(X_train_50) > 0:\n                model3 = self.trainer.train_model(X_train_50, y_train_50, X_valid, y_valid, weights_50)\n                \n                valid_idx_recent = valid_idx[valid_idx >= cutoff_50]\n                if len(valid_idx_recent) > 0:\n                    X_valid_recent = train_df.iloc[valid_idx_recent][features]\n                    oof_predictions['model_3'][valid_idx_recent] = model3.predict(X_valid_recent)\n                \n                valid_idx_old = valid_idx[valid_idx < cutoff_50]\n                if len(valid_idx_old) > 0:\n                    oof_predictions['model_3'][valid_idx_old] = oof_predictions['model_1'][valid_idx_old]\n                \n                test_predictions['model_3'] += model3.predict(X_test) / self.config.n_folds\n        \n        return {\n            'oof_predictions': oof_predictions,\n            'test_predictions': test_predictions,\n            'train_labels': train_df[\"label\"]\n        }\n    \n    def save_feature_importance(self, train_df: pd.DataFrame, test_df: pd.DataFrame):\n        \"\"\"Generate and save SHAP feature importance analysis\"\"\"\n        try:\n            import shap\n            \n            features = [c for c in train_df.columns if c != \"label\"]\n            sample_weights = self.data_processor.create_time_weights(\n                len(train_df), self.config.decay_factor\n            )\n            \n            print(\"\\nGenerating SHAP feature importance...\")\n            \n            # Train model for SHAP analysis\n            model = xgb.XGBRegressor(**self.config.xgb_params)\n            model.fit(\n                train_df[features], \n                train_df[\"label\"],\n                sample_weight=sample_weights,\n                verbose=0\n            )\n            \n            # Calculate SHAP values on a sample\n            sample_size = min(1000, len(test_df))\n            test_sample = test_df[features].iloc[:sample_size]\n            \n            explainer = shap.TreeExplainer(model)\n            shap_values = explainer.shap_values(test_sample)\n            \n            # Create feature importance dataframe\n            feature_importance = pd.DataFrame({\n                'feature': features,\n                'importance': np.abs(shap_values).mean(axis=0)\n            }).sort_values('importance', ascending=False)\n            \n            feature_importance.to_csv(self.config.shap_features_path, index=False)\n            print(f\"Feature importance saved to {self.config.shap_features_path}\")\n            \n            # Register SHAP features output with global configuration\n            global_config.register_model_output(\n                self.config.model_name,\n                'shap_features',\n                self.config.shap_features_path,\n                'analysis',\n                metadata={'top_features': feature_importance.head(10)['feature'].tolist()}\n            )\n            \n            # Save top features\n            print(\"\\nTop 10 most important features:\")\n            for idx, row in feature_importance.head(10).iterrows():\n                print(f\"  {row['feature']}: {row['importance']:.4f}\")\n            \n        except Exception as e:\n            print(f\"Warning: Could not generate SHAP feature importance: {e}\")\n    \n    def run(self) -> float:\n        \"\"\"Execute the complete XGBoost three-model pipeline\"\"\"\n        print(\"\\nStarting XGBoost Three-Model Pipeline\")\n        print(\"=\"*80)\n        \n        # Update model status\n        global_config.update_model_status(self.config.model_name, \"running\")\n        \n        # Load data\n        train_df, test_df, sample_df = self.data_processor.load_data(self.config)\n        \n        # Run cross-validation\n        print(\"\\nRunning cross-validation...\")\n        cv_results = self.run_cross_validation(train_df, test_df)\n        \n        # Evaluate models\n        print(\"\\nEvaluating model performance...\")\n        ensemble_results = self.ensemble_builder.evaluate_models(\n            cv_results['oof_predictions'],\n            cv_results['train_labels']\n        )\n        \n        # Display results\n        print(\"\\nModel Performance Summary:\")\n        for _, row in ensemble_results.iterrows():\n            if pd.notna(row['pearson_correlation']):\n                print(f\"  {row['model']}: {row['pearson_correlation']:.4f}\")\n        \n        # Get best ensemble score\n        best_score = ensemble_results['pearson_correlation'].max()\n        best_model = ensemble_results.loc[\n            ensemble_results['pearson_correlation'].idxmax(), 'model'\n        ]\n        \n        print(f\"\\nBest performing approach: {best_model} (score: {best_score:.4f})\")\n        \n        # Create final predictions\n        oof_scores = {\n            name: ensemble_results[ensemble_results['model'] == name]['pearson_correlation'].values[0]\n            for name in cv_results['test_predictions'].keys()\n        }\n        \n        final_predictions = self.ensemble_builder.create_final_ensemble(\n            cv_results['test_predictions'],\n            oof_scores\n        )\n        \n        # Save submission\n        sample_df[\"prediction\"] = final_predictions\n        sample_df.to_csv(self.config.submission_file, index=False)\n        print(f\"\\nSubmission saved to {self.config.submission_file}\")\n        \n        # Register submission with global configuration\n        global_config.register_model_output(\n            self.config.model_name,\n            'submission',\n            self.config.submission_file,\n            'submission',\n            metadata={\n                'best_score': float(best_score),\n                'best_model': best_model,\n                'n_models': len(self.config.model_configs)\n            }\n        )\n        \n        # Save ensemble results\n        ensemble_results.to_csv(self.config.results_file, index=False)\n        \n        # Register ensemble results with global configuration\n        global_config.register_model_output(\n            self.config.model_name,\n            'ensemble_results',\n            self.config.results_file,\n            'analysis',\n            metadata={\n                'model_scores': ensemble_results.set_index('model')['pearson_correlation'].to_dict()\n            }\n        )\n        \n        # Generate feature importance\n        self.save_feature_importance(train_df, test_df)\n        \n        print(\"\\nXGBoost pipeline completed successfully\")\n        \n        return best_score\n\n# Main execution\nif __name__ == \"__main__\":\n    try:\n        print(\"\\n📊 Running XGBoost Three-Model Pipeline\")\n        print(\"-\"*80)\n        \n        # Clean memory before starting\n        aggressive_memory_cleanup()\n        \n        # Create and run pipeline\n        pipeline = XGBoostThreeModelPipeline()\n        final_score = pipeline.run()\n        \n        # Update global configuration with success\n        global_config.update_model_status('xgboost', 'completed', score=final_score)\n        print(\"\\n✅ XGBoost pipeline completed successfully\")\n        \n    except Exception as e:\n        # Update global configuration with failure\n        error_msg = str(e)\n        global_config.update_model_status('xgboost', 'failed', error_message=error_msg)\n        print(f\"\\n❌ XGBoost pipeline failed: {error_msg}\")\n        raise\n        \n    finally:\n        # Clean up memory\n        aggressive_memory_cleanup()\n        \n        # Display current execution status\n        print(\"\\n\" + \"=\"*80)\n        print(\"Current Execution Status:\")\n        print(global_config.get_execution_summary())\n        print(\"=\"*80)","metadata":{"trusted":true,"execution":{"iopub.status.busy":"2025-06-04T16:05:34.003635Z","iopub.execute_input":"2025-06-04T16:05:34.004006Z","iopub.status.idle":"2025-06-04T16:09:38.791145Z","shell.execute_reply.started":"2025-06-04T16:05:34.003986Z","shell.execute_reply":"2025-06-04T16:09:38.790260Z"},"jupyter":{"outputs_hidden":true,"source_hidden":true},"collapsed":true},"outputs":[],"execution_count":null},{"cell_type":"code","source":"# AUTOENCODER DEEP MLP PIPELINE IMPLEMENTATION\n# !/usr/bin/env python\n# -*- coding: utf-8 -*-\n\"\"\"\nDRW Crypto Market Prediction - AutoEncoder Deep MLP Pipeline\nThis module implements a deep learning approach using autoencoders for feature extraction\ncombined with a multi-layer perceptron with residual connections for prediction\n\"\"\"\n\nimport subprocess\nimport sys\nimport os\nimport gc\nimport warnings\nimport json\nimport random\nimport pandas as pd\nimport numpy as np\nfrom scipy.stats import pearsonr\nfrom typing import List, Dict, Tuple, Optional, Any\nfrom pathlib import Path\nfrom dataclasses import dataclass, field\nfrom sklearn.model_selection import TimeSeriesSplit\nfrom sklearn.linear_model import Ridge\nfrom sklearn.preprocessing import StandardScaler\nimport matplotlib.pyplot as plt\n\nwarnings.filterwarnings('ignore')\n\n# Install required packages for this pipeline\nprint(\"Installing packages for AutoEncoder Deep MLP pipeline...\")\npackages_to_install = [\n    'torch==2.0.1',\n    'torchvision==0.15.2',\n    'tqdm==4.65.0'\n]\n\n# Special handling for PyTorch to ensure CPU-only installation\nprint(\"Installing PyTorch (CPU version)...\")\nsubprocess.check_call([\n    sys.executable, \"-m\", \"pip\", \"install\", \n    \"torch==2.0.1\", \"torchvision==0.15.2\",\n    \"--index-url\", \"https://download.pytorch.org/whl/cpu\",\n    \"--quiet\"\n])\n\n# Install remaining packages\nsubprocess.check_call([sys.executable, \"-m\", \"pip\", \"install\", \"tqdm==4.65.0\", \"--quiet\"])\n\n@dataclass\nclass AutoEncoderConfiguration:\n    \"\"\"Configuration for AutoEncoder Deep MLP Pipeline\"\"\"\n    \n    # Model identification\n    model_name: str = \"autoencoder\"\n    model_directory: str = \"\"\n    \n    # Model-specific features\n    feature_columns: List[str] = field(default_factory=list)\n    \n    # AutoEncoder architecture parameters\n    encoding_size: int = 128\n    hidden_size: int = 256\n    dropout: float = 0.7\n    ae_dropout: float = 0.3\n    noise_std: float = 0.05\n    num_blocks: int = 8\n    \n    # Training parameters\n    num_epochs: int = 80\n    batch_size: int = 4096\n    learning_rate: float = 0.0001\n    weight_decay: float = 5e-3\n    patience: int = 10\n    min_lr: float = 1e-6\n    \n    # Loss weights\n    mse_weight: float = 0.25\n    corr_weight: float = 0.6\n    ae_weight: float = 0.15\n    \n    # Cross-validation parameters\n    n_splits: int = 5\n    max_train_size: int = 100_000_000\n    gap: int = 1\n    \n    # Random seed\n    seed: int = 42\n    \n    def __post_init__(self):\n        \"\"\"Initialize configuration and register models\"\"\"\n        if not self.model_directory:\n            self.model_directory = os.path.join(global_config.base_dir, \"autoencoder_deepmlp\")\n        \n        # Register both autoencoder variants with global configuration\n        global_config.register_model(\"autoencoder_simple\", self.model_directory)\n        global_config.register_model(\"autoencoder_weighted\", self.model_directory)\n        \n        # Set default features if not provided\n        if not self.feature_columns:\n            self.feature_columns = [\n                \"X863\", \"X856\", \"X344\", \"X598\", \"X862\", \"X385\", \"X852\", \"X603\", \"X860\", \"X674\",\n                \"X415\", \"X345\", \"X137\", \"X855\", \"X174\", \"X302\", \"X178\", \"X532\", \"X168\", \"X612\",\n                \"bid_qty\", \"ask_qty\", \"buy_qty\", \"sell_qty\", \"volume\",\n                \"X188\", \"X207\", \"X219\", \"X233\", \"X245\"\n            ]\n        \n        # Output paths\n        self.intermediate_dir = os.path.join(self.model_directory, \"fold_models\")\n        self.simple_submission_path = os.path.join(self.model_directory, \"ensemble_simple_submission.csv\")\n        self.weighted_submission_path = os.path.join(self.model_directory, \"ensemble_weighted_submission.csv\")\n        self.config_path = os.path.join(self.model_directory, \"config.json\")\n        self.performance_metrics_path = os.path.join(self.model_directory, \"performance_metrics.json\")\n        \n        # Ensure directories exist\n        Path(self.model_directory).mkdir(parents=True, exist_ok=True)\n        Path(self.intermediate_dir).mkdir(parents=True, exist_ok=True)\n    \n    def save(self):\n        \"\"\"Save configuration to JSON file\"\"\"\n        config_dict = {k: v for k, v in self.__dict__.items() \n                      if not k.startswith('_') and k not in ['model_directory', 'intermediate_dir']}\n        with open(self.config_path, 'w') as f:\n            json.dump(config_dict, f, indent=2)\n\nclass RandomSeedManager:\n    \"\"\"Manages random seeds across all libraries for reproducibility\"\"\"\n    \n    @staticmethod\n    def set_all_seeds(seed: int):\n        \"\"\"Set random seeds for all relevant libraries\"\"\"\n        random.seed(seed)\n        np.random.seed(seed)\n        os.environ['PYTHONHASHSEED'] = str(seed)\n        \n        if 'torch' in sys.modules:\n            import torch\n            torch.manual_seed(seed)\n            torch.cuda.manual_seed(seed)\n            torch.cuda.manual_seed_all(seed)\n            torch.backends.cudnn.deterministic = True\n            torch.backends.cudnn.benchmark = False\n\nclass TorchEnvironmentManager:\n    \"\"\"Manages PyTorch environment configuration to avoid conflicts\"\"\"\n    \n    @staticmethod\n    def configure_environment():\n        \"\"\"Configure PyTorch environment for CPU execution\"\"\"\n        # Force CPU usage to avoid GPU/CUDA issues\n        os.environ['CUDA_VISIBLE_DEVICES'] = ''\n        os.environ['PYTORCH_CUDA_ALLOC_CONF'] = 'max_split_size_mb:128'\n        os.environ['PYTORCH_NO_CUDA_MEMORY_CACHING'] = '1'\n        \n        # Disable Triton to avoid registration conflicts\n        os.environ['TRITON_CACHE_DIR'] = f'/tmp/triton_cache_{os.getpid()}'\n        os.environ['DISABLE_TRITON'] = '1'\n        \n        # Clean torch modules if already loaded\n        torch_modules = [\n            'torch', 'torchvision', 'torchaudio', 'triton',\n            'torch.nn', 'torch.optim', 'torch.utils', 'torch.cuda'\n        ]\n        \n        for module in list(sys.modules.keys()):\n            if any(module.startswith(torch_mod) for torch_mod in torch_modules):\n                try:\n                    del sys.modules[module]\n                except:\n                    pass\n        \n        gc.collect()\n\nclass AutoEncoderDataProcessor:\n    \"\"\"Data processing utilities for AutoEncoder pipeline\"\"\"\n    \n    @staticmethod\n    def load_and_prepare_data(config: AutoEncoderConfiguration) -> Tuple[np.ndarray, np.ndarray, np.ndarray, pd.DataFrame]:\n        \"\"\"Load and prepare data for AutoEncoder training\"\"\"\n        print(\"Loading data...\")\n        train_df = pd.read_parquet(global_config.train_path)\n        test_df = pd.read_parquet(global_config.test_path)\n        sample_submission = pd.read_csv(global_config.sample_sub_path)\n        \n        # Verify feature availability\n        available_features = [col for col in config.feature_columns if col in train_df.columns]\n        missing_features = set(config.feature_columns) - set(available_features)\n        \n        if missing_features:\n            print(f\"Warning: {len(missing_features)} features not found in data\")\n            print(f\"Missing features: {missing_features}\")\n        \n        # Select available features\n        train_df = train_df[available_features + ['label']]\n        test_df = test_df[available_features]\n        \n        # Optimize memory usage\n        train_df = reduce_memory_usage(train_df, verbose=False)\n        test_df = reduce_memory_usage(test_df, verbose=False)\n        \n        # Extract arrays\n        X_full = train_df[available_features].values\n        y_full = train_df['label'].values\n        X_test = test_df[available_features].values\n        \n        # Handle invalid values\n        X_full = np.nan_to_num(X_full, nan=0.0, posinf=0.0, neginf=0.0)\n        X_test = np.nan_to_num(X_test, nan=0.0, posinf=0.0, neginf=0.0)\n        \n        print(f\"Data loaded - Train: {X_full.shape}, Test: {X_test.shape}\")\n        print(f\"Using {len(available_features)} features\")\n        \n        return X_full, y_full, X_test, sample_submission\n\nclass AutoEncoderNeuralNetwork:\n    \"\"\"Neural network components for AutoEncoder architecture\"\"\"\n    \n    def __init__(self, config: AutoEncoderConfiguration):\n        self.config = config\n        self.device = None\n        self.torch = None\n        self.nn = None\n        self.F = None\n    \n    def initialize_pytorch(self):\n        \"\"\"Initialize PyTorch modules after environment configuration\"\"\"\n        try:\n            import torch\n            import torch.nn as nn\n            import torch.nn.functional as F\n            \n            self.torch = torch\n            self.nn = nn\n            self.F = F\n            self.device = torch.device('cpu')\n            \n            print(f\"PyTorch version: {torch.__version__}\")\n            print(f\"Using device: {self.device}\")\n            \n            return True\n            \n        except Exception as e:\n            print(f\"Error initializing PyTorch: {e}\")\n            return False\n    \n    def create_autoencoder(self, input_size: int):\n        \"\"\"Create autoencoder neural network\"\"\"\n        \n        class Swish(self.nn.Module):\n            def forward(self, x):\n                return x * self.torch.sigmoid(x)\n        \n        class GaussianNoise(self.nn.Module):\n            def __init__(self, std: float = 0.05):\n                super().__init__()\n                self.std = std\n                \n            def forward(self, x):\n                if self.training:\n                    noise = self.torch.randn_like(x) * self.std\n                    return x + noise\n                return x\n        \n        class AutoEncoder(self.nn.Module):\n            def __init__(self, input_size: int, encoding_size: int, dropout: float):\n                super().__init__()\n                \n                # Encoder layers\n                self.encoder = self.nn.Sequential(\n                    self.nn.Linear(input_size, input_size // 2),\n                    self.nn.BatchNorm1d(input_size // 2),\n                    Swish(),\n                    self.nn.Dropout(dropout),\n                    \n                    self.nn.Linear(input_size // 2, input_size // 4),\n                    self.nn.BatchNorm1d(input_size // 4),\n                    Swish(),\n                    self.nn.Dropout(dropout),\n                    \n                    self.nn.Linear(input_size // 4, encoding_size),\n                    self.nn.BatchNorm1d(encoding_size),\n                    Swish()\n                )\n                \n                # Decoder layers\n                self.decoder = self.nn.Sequential(\n                    self.nn.Linear(encoding_size, input_size // 4),\n                    self.nn.BatchNorm1d(input_size // 4),\n                    Swish(),\n                    self.nn.Dropout(dropout),\n                    \n                    self.nn.Linear(input_size // 4, input_size // 2),\n                    self.nn.BatchNorm1d(input_size // 2),\n                    Swish(),\n                    self.nn.Dropout(dropout),\n                    \n                    self.nn.Linear(input_size // 2, input_size)\n                )\n                \n            def forward(self, x):\n                encoded = self.encoder(x)\n                decoded = self.decoder(encoded)\n                return encoded, decoded\n        \n        return AutoEncoder(input_size, self.config.encoding_size, self.config.ae_dropout)\n    \n    def create_full_model(self, input_size: int):\n        \"\"\"Create complete model with autoencoder and prediction head\"\"\"\n        \n        autoencoder = self.create_autoencoder(input_size)\n        config = self.config\n        \n        class CryptoMLPWithAutoEncoder(self.nn.Module):\n            def __init__(self):\n                super().__init__()\n                \n                self.noise_layer = self.nn.Sequential()  # Will add GaussianNoise\n                self.autoencoder = autoencoder\n                \n                combined_input_size = input_size + config.encoding_size\n                self.input_bn = self.nn.BatchNorm1d(combined_input_size)\n                \n                # Initial block\n                self.initial_block = self.nn.Sequential(\n                    self.nn.Linear(combined_input_size, config.hidden_size),\n                    self.nn.BatchNorm1d(config.hidden_size),\n                    self.nn.ReLU(),\n                    self.nn.Dropout(config.dropout),\n                    \n                    self.nn.Linear(config.hidden_size, config.hidden_size),\n                    self.nn.BatchNorm1d(config.hidden_size),\n                    self.nn.ReLU(),\n                    self.nn.Dropout(config.dropout),\n                    \n                    self.nn.Linear(config.hidden_size, config.hidden_size),\n                    self.nn.BatchNorm1d(config.hidden_size),\n                    self.nn.ReLU()\n                )\n                \n                # Residual blocks\n                self.residual_blocks = self.nn.ModuleList()\n                for _ in range(config.num_blocks):\n                    block = self.nn.Sequential(\n                        self.nn.Linear(config.hidden_size, config.hidden_size),\n                        self.nn.BatchNorm1d(config.hidden_size),\n                        self.nn.ReLU(),\n                        self.nn.Dropout(config.dropout),\n                        \n                        self.nn.Linear(config.hidden_size, config.hidden_size),\n                        self.nn.BatchNorm1d(config.hidden_size),\n                        self.nn.ReLU()\n                    )\n                    self.residual_blocks.append(block)\n                \n                # Output layer\n                self.output = self.nn.Linear(config.hidden_size, 1)\n                \n            def forward(self, x, return_ae_loss=False):\n                # Add noise during training\n                if self.training:\n                    x = x + self.torch.randn_like(x) * config.noise_std\n                \n                # Autoencoder\n                encoded, decoded = self.autoencoder(x)\n                \n                # Combine original and encoded features\n                x_combined = self.torch.cat([x, encoded], dim=1)\n                x_combined = self.input_bn(x_combined)\n                \n                # Initial transformation\n                x_hidden = self.initial_block(x_combined)\n                \n                # Residual connections\n                for block in self.residual_blocks:\n                    x_residual = block(x_hidden)\n                    x_hidden = x_hidden + x_residual\n                \n                # Output\n                output = self.output(x_hidden)\n                \n                if return_ae_loss:\n                    return output, decoded, x\n                return output\n        \n        return CryptoMLPWithAutoEncoder()\n\nclass AutoEncoderTrainer:\n    \"\"\"Training logic for AutoEncoder models\"\"\"\n    \n    def __init__(self, config: AutoEncoderConfiguration, nn_builder: AutoEncoderNeuralNetwork):\n        self.config = config\n        self.nn_builder = nn_builder\n    \n    def train_fold(self, X_train: np.ndarray, y_train: np.ndarray,\n                   X_val: np.ndarray, y_val: np.ndarray,\n                   X_test: np.ndarray, fold_idx: int) -> Dict[str, Any]:\n        \"\"\"Train a single fold of the AutoEncoder model\"\"\"\n        \n        torch = self.nn_builder.torch\n        nn = self.nn_builder.nn\n        device = self.nn_builder.device\n        \n        # Create data loaders\n        from torch.utils.data import Dataset, DataLoader\n        \n        class CryptoDataset(Dataset):\n            def __init__(self, features: np.ndarray, labels: Optional[np.ndarray] = None):\n                self.features = torch.FloatTensor(features)\n                self.labels = torch.FloatTensor(labels) if labels is not None else None\n                \n            def __len__(self):\n                return len(self.features)\n                \n            def __getitem__(self, idx):\n                if self.labels is not None:\n                    return self.features[idx], self.labels[idx]\n                return self.features[idx]\n        \n        train_dataset = CryptoDataset(X_train, y_train)\n        val_dataset = CryptoDataset(X_val, y_val)\n        test_dataset = CryptoDataset(X_test)\n        \n        train_loader = DataLoader(train_dataset, batch_size=self.config.batch_size, \n                                shuffle=True, num_workers=0)\n        val_loader = DataLoader(val_dataset, batch_size=self.config.batch_size, \n                              shuffle=False, num_workers=0)\n        test_loader = DataLoader(test_dataset, batch_size=self.config.batch_size, \n                               shuffle=False, num_workers=0)\n        \n        # Create model\n        num_features = X_train.shape[1]\n        model = self.nn_builder.create_full_model(num_features).to(device)\n        \n        # Initialize weights\n        def init_weights(m):\n            if isinstance(m, nn.Linear):\n                torch.nn.init.xavier_uniform_(m.weight)\n                m.bias.data.fill_(0.01)\n        \n        model.apply(init_weights)\n        \n        # Loss function\n        class CombinedLoss(nn.Module):\n            def __init__(self):\n                super().__init__()\n                self.mse = nn.MSELoss()\n                \n            def forward(self, y_pred, y_true, decoded=None, original=None):\n                # Main prediction loss\n                mse_loss = self.mse(y_pred, y_true)\n                \n                # Correlation loss\n                y_pred_centered = y_pred - y_pred.mean()\n                y_true_centered = y_true - y_true.mean()\n                \n                correlation = torch.sum(y_pred_centered * y_true_centered) / (\n                    torch.sqrt(torch.sum(y_pred_centered ** 2)) * \n                    torch.sqrt(torch.sum(y_true_centered ** 2)) + 1e-8\n                )\n                \n                # Autoencoder reconstruction loss\n                ae_loss = self.mse(decoded, original) if decoded is not None else 0\n                \n                total_loss = (self.config.mse_weight * mse_loss - \n                            self.config.corr_weight * correlation + \n                            self.config.ae_weight * ae_loss)\n                \n                return total_loss, correlation.item()\n        \n        criterion = CombinedLoss()\n        optimizer = torch.optim.Adam(model.parameters(), lr=self.config.learning_rate,\n                                   weight_decay=self.config.weight_decay)\n        scheduler = torch.optim.lr_scheduler.ReduceLROnPlateau(optimizer, mode='max',\n                                                              patience=5, factor=0.5,\n                                                              min_lr=self.config.min_lr)\n        \n        # Training loop\n        best_val_corr = -float('inf')\n        best_model_state = None\n        early_stop_counter = 0\n        \n        from tqdm import tqdm\n        \n        for epoch in range(self.config.num_epochs):\n            # Training\n            model.train()\n            train_losses = []\n            \n            for batch_features, batch_labels in tqdm(train_loader, desc=f\"Epoch {epoch+1}\", leave=False):\n                batch_features = batch_features.to(device)\n                batch_labels = batch_labels.to(device)\n                \n                optimizer.zero_grad()\n                \n                outputs, decoded, original = model(batch_features, return_ae_loss=True)\n                loss, _ = criterion(outputs, batch_labels, decoded, original)\n                \n                loss.backward()\n                torch.nn.utils.clip_grad_norm_(model.parameters(), max_norm=1.0)\n                optimizer.step()\n                \n                train_losses.append(loss.item())\n            \n            # Validation\n            model.eval()\n            val_predictions = []\n            val_targets = []\n            \n            with torch.no_grad():\n                for batch_features, batch_labels in val_loader:\n                    batch_features = batch_features.to(device)\n                    outputs = model(batch_features)\n                    \n                    val_predictions.extend(outputs.cpu().numpy().flatten())\n                    val_targets.extend(batch_labels.numpy().flatten())\n            \n            # Calculate validation correlation\n            val_corr = pearsonr(val_predictions, val_targets)[0]\n            scheduler.step(val_corr)\n            \n            # Early stopping check\n            if val_corr > best_val_corr:\n                best_val_corr = val_corr\n                best_model_state = model.state_dict().copy()\n                early_stop_counter = 0\n            else:\n                early_stop_counter += 1\n            \n            if (epoch + 1) % 10 == 0:\n                print(f\"  Epoch {epoch+1}/{self.config.num_epochs} | Val Corr: {val_corr:.6f}\")\n            \n            if early_stop_counter >= self.config.patience:\n                print(f\"  Early stopping at epoch {epoch+1}\")\n                break\n        \n        # Load best model and make predictions\n        model.load_state_dict(best_model_state)\n        model.eval()\n        \n        test_predictions = []\n        with torch.no_grad():\n            for batch_features in tqdm(test_loader, desc=\"Predicting\", leave=False):\n                batch_features = batch_features.to(device)\n                outputs = model(batch_features)\n                test_predictions.append(outputs.cpu().numpy())\n        \n        test_predictions = np.vstack(test_predictions).flatten()\n        \n        # Save model state\n        model_path = os.path.join(self.config.intermediate_dir, f'fold_{fold_idx}_model.pt')\n        torch.save(best_model_state, model_path)\n        \n        return {\n            'fold_idx': fold_idx,\n            'best_val_corr': best_val_corr,\n            'test_predictions': test_predictions\n        }\n\nclass AutoEncoderPipeline:\n    \"\"\"Main pipeline orchestrating the AutoEncoder Deep MLP approach\"\"\"\n    \n    def __init__(self, config: Optional[AutoEncoderConfiguration] = None):\n        self.config = config or AutoEncoderConfiguration()\n        self.data_processor = AutoEncoderDataProcessor()\n    \n    def create_fallback_predictions(self, X_full: np.ndarray, y_full: np.ndarray, \n                                  X_test: np.ndarray, sample_submission: pd.DataFrame) -> Dict[str, Any]:\n        \"\"\"Create predictions using Ridge regression as fallback\"\"\"\n        print(\"\\nUsing Ridge regression fallback...\")\n        \n        scaler = StandardScaler()\n        X_train_scaled = scaler.fit_transform(X_full)\n        X_test_scaled = scaler.transform(X_test)\n        \n        model = Ridge(alpha=1.0, random_state=self.config.seed)\n        model.fit(X_train_scaled, y_full)\n        \n        predictions = model.predict(X_test_scaled)\n        \n        # Estimate score\n        from sklearn.model_selection import cross_val_score\n        scores = cross_val_score(model, X_train_scaled, y_full, cv=3,\n                               scoring=lambda est, X, y: pearsonr(est.predict(X), y)[0])\n        final_score = np.mean(scores)\n        \n        return {\n            'predictions': predictions,\n            'score': final_score,\n            'method': 'ridge_fallback'\n        }\n    \n    def save_performance_metrics(self, all_results: List[Dict], weights: np.ndarray, final_score: float):\n        \"\"\"Save detailed performance metrics\"\"\"\n        performance_data = {\n            'final_score': float(final_score),\n            'method': 'neural_network' if all_results else 'ridge_fallback',\n            'n_folds': len(all_results),\n            'fold_results': []\n        }\n        \n        if all_results:\n            for i, result in enumerate(all_results):\n                performance_data['fold_results'].append({\n                    'fold': result['fold_idx'],\n                    'validation_correlation': float(result['best_val_corr']),\n                    'weight': float(weights[i]) if i < len(weights) else 0.0\n                })\n            \n            correlations = [r['best_val_corr'] for r in all_results]\n            performance_data['statistics'] = {\n                'mean_correlation': float(np.mean(correlations)),\n                'std_correlation': float(np.std(correlations)),\n                'min_correlation': float(np.min(correlations)),\n                'max_correlation': float(np.max(correlations))\n            }\n        \n        with open(self.config.performance_metrics_path, 'w') as f:\n            json.dump(performance_data, f, indent=2)\n        \n        # Register performance metrics with global configuration\n        global_config.register_model_output(\n            'autoencoder_simple',\n            'performance_metrics',\n            self.config.performance_metrics_path,\n            'analysis',\n            metadata=performance_data\n        )\n    \n    def run(self) -> float:\n        \"\"\"Execute the complete AutoEncoder pipeline\"\"\"\n        print(\"\\nStarting AutoEncoder Deep MLP Pipeline\")\n        print(\"=\"*80)\n        \n        # Update status for both model variants\n        global_config.update_model_status('autoencoder_simple', 'running')\n        global_config.update_model_status('autoencoder_weighted', 'running')\n        \n        # Configure PyTorch environment\n        print(\"Configuring PyTorch environment...\")\n        TorchEnvironmentManager.configure_environment()\n        \n        # Load data\n        X_full, y_full, X_test, sample_submission = self.data_processor.load_and_prepare_data(self.config)\n        \n        # Initialize neural network builder\n        nn_builder = AutoEncoderNeuralNetwork(self.config)\n        \n        if not nn_builder.initialize_pytorch():\n            print(\"PyTorch initialization failed, using fallback method\")\n            fallback_results = self.create_fallback_predictions(X_full, y_full, X_test, sample_submission)\n            \n            # Save fallback predictions\n            submission = sample_submission.copy()\n            submission.iloc[:, 1] = fallback_results['predictions']\n            submission.to_csv(self.config.simple_submission_path, index=False)\n            submission.to_csv(self.config.weighted_submission_path, index=False)\n            \n            # Register outputs with global configuration\n            global_config.register_model_output(\n                'autoencoder_simple',\n                'ensemble_simple_submission',\n                self.config.simple_submission_path,\n                'submission',\n                metadata={'method': 'ridge_fallback', 'score': fallback_results['score']}\n            )\n            \n            global_config.register_model_output(\n                'autoencoder_weighted',\n                'ensemble_weighted_submission',\n                self.config.weighted_submission_path,\n                'submission',\n                metadata={'method': 'ridge_fallback', 'score': fallback_results['score']}\n            )\n            \n            return fallback_results['score']\n        \n        # Create trainer\n        trainer = AutoEncoderTrainer(self.config, nn_builder)\n        \n        # Time series cross-validation\n        tss = TimeSeriesSplit(n_splits=self.config.n_splits,\n                            max_train_size=self.config.max_train_size,\n                            gap=self.config.gap)\n        \n        all_results = []\n        all_test_predictions = []\n        \n        for fold_idx, (train_idx, val_idx) in enumerate(tss.split(X_full)):\n            print(f\"\\nTraining Fold {fold_idx + 1}/{self.config.n_splits}\")\n            \n            RandomSeedManager.set_all_seeds(self.config.seed + fold_idx)\n            \n            X_train = X_full[train_idx]\n            X_val = X_full[val_idx]\n            y_train = y_full[train_idx]\n            y_val = y_full[val_idx]\n            \n            print(f\"  Train samples: {len(X_train):,}\")\n            print(f\"  Val samples: {len(X_val):,}\")\n            \n            try:\n                fold_results = trainer.train_fold(X_train, y_train, X_val, y_val, X_test, fold_idx)\n                all_results.append(fold_results)\n                all_test_predictions.append(fold_results['test_predictions'])\n                \n                print(f\"  Fold {fold_idx + 1} completed | Best Corr: {fold_results['best_val_corr']:.6f}\")\n                \n            except Exception as e:\n                print(f\"  Error in fold {fold_idx + 1}: {e}\")\n                print(\"  Continuing with remaining folds...\")\n            \n            # Clean up memory\n            aggressive_memory_cleanup()\n        \n        if not all_results:\n            print(\"All folds failed, using fallback method\")\n            fallback_results = self.create_fallback_predictions(X_full, y_full, X_test, sample_submission)\n            \n            submission = sample_submission.copy()\n            submission.iloc[:, 1] = fallback_results['predictions']\n            submission.to_csv(self.config.simple_submission_path, index=False)\n            submission.to_csv(self.config.weighted_submission_path, index=False)\n            \n            # Register outputs with global configuration\n            global_config.register_model_output(\n                'autoencoder_simple',\n                'ensemble_simple_submission',\n                self.config.simple_submission_path,\n                'submission',\n                metadata={'method': 'ridge_fallback', 'score': fallback_results['score']}\n            )\n            \n            global_config.register_model_output(\n                'autoencoder_weighted',\n                'ensemble_weighted_submission',\n                self.config.weighted_submission_path,\n                'submission',\n                metadata={'method': 'ridge_fallback', 'score': fallback_results['score']}\n            )\n            \n            return fallback_results['score']\n        \n        # Create ensemble predictions\n        print(\"\\nCreating ensemble predictions...\")\n        \n        # Simple average ensemble\n        ensemble_predictions = np.mean(all_test_predictions, axis=0)\n        \n        # Weighted ensemble based on validation scores\n        weights = np.array([r['best_val_corr'] for r in all_results])\n        weights = np.maximum(weights, 0)\n        weights = weights / weights.sum()\n        \n        weighted_ensemble_predictions = np.average(all_test_predictions, axis=0, weights=weights)\n        \n        # Save predictions\n        submission = sample_submission.copy()\n        submission.iloc[:, 1] = ensemble_predictions\n        submission.to_csv(self.config.simple_submission_path, index=False)\n        \n        # Register simple ensemble output\n        global_config.register_model_output(\n            'autoencoder_simple',\n            'ensemble_simple_submission',\n            self.config.simple_submission_path,\n            'submission',\n            metadata={\n                'method': 'simple_average',\n                'n_folds': len(all_results),\n                'score': float(np.mean([r['best_val_corr'] for r in all_results]))\n            }\n        )\n        \n        submission = sample_submission.copy()\n        submission.iloc[:, 1] = weighted_ensemble_predictions\n        submission.to_csv(self.config.weighted_submission_path, index=False)\n        \n        # Register weighted ensemble output\n        global_config.register_model_output(\n            'autoencoder_weighted',\n            'ensemble_weighted_submission',\n            self.config.weighted_submission_path,\n            'submission',\n            metadata={\n                'method': 'weighted_average',\n                'n_folds': len(all_results),\n                'weights': weights.tolist(),\n                'score': float(np.mean([r['best_val_corr'] for r in all_results]))\n            }\n        )\n        \n        # Save configuration\n        self.config.save()\n        \n        # Register configuration\n        global_config.register_model_output(\n            'autoencoder_simple',\n            'config',\n            self.config.config_path,\n            'config',\n            metadata={'encoding_size': self.config.encoding_size, 'hidden_size': self.config.hidden_size}\n        )\n        \n        # Calculate final score\n        correlations = [r['best_val_corr'] for r in all_results]\n        final_score = np.mean(correlations)\n        \n        # Save performance metrics\n        self.save_performance_metrics(all_results, weights, final_score)\n        \n        print(f\"\\nAutoEncoder Ensemble Statistics:\")\n        print(f\"  Average Correlation: {final_score:.6f}\")\n        print(f\"  Best Single Fold: {np.max(correlations):.6f}\")\n        print(f\"  Fold Weights: {weights}\")\n        \n        print(\"\\nAutoEncoder pipeline completed successfully\")\n        \n        return final_score\n\n# Main execution\nif __name__ == \"__main__\":\n    try:\n        print(\"\\n🧠 Running AutoEncoder Deep MLP Pipeline\")\n        print(\"-\"*80)\n        \n        # Clean memory before starting\n        aggressive_memory_cleanup()\n        \n        # Create configuration\n        config = AutoEncoderConfiguration(\n            num_epochs=80,\n            encoding_size=128,\n            hidden_size=256,\n            num_blocks=8,\n            dropout=0.7,\n            noise_std=0.05\n        )\n        \n        # Create and run pipeline\n        pipeline = AutoEncoderPipeline(config)\n        final_score = pipeline.run()\n        \n        # Update global configuration with success\n        global_config.update_model_status('autoencoder_simple', 'completed', score=final_score)\n        global_config.update_model_status('autoencoder_weighted', 'completed', score=final_score)\n        print(\"\\n✅ AutoEncoder pipeline completed successfully\")\n        \n    except Exception as e:\n        # Update global configuration with failure\n        error_msg = str(e)\n        global_config.update_model_status('autoencoder_simple', 'failed', error_message=error_msg)\n        global_config.update_model_status('autoencoder_weighted', 'failed', error_message=error_msg)\n        print(f\"\\n❌ AutoEncoder pipeline failed: {error_msg}\")\n        raise\n        \n    finally:\n        # Clean up memory\n        aggressive_memory_cleanup()\n        \n        # Display current execution status\n        print(\"\\n\" + \"=\"*80)\n        print(\"Current Execution Status:\")\n        print(global_config.get_execution_summary())\n        print(\"=\"*80)","metadata":{"trusted":true,"jupyter":{"source_hidden":true}},"outputs":[],"execution_count":null},{"cell_type":"code","source":"# MARKET MICROSTRUCTURE XGBOOST PIPELINE IMPLEMENTATION\n# ===========================================================================================\n# ===========================================================================================\n\n#!/usr/bin/env python\n# -*- coding: utf-8 -*-\n\"\"\"\nDRW Crypto Market Prediction - Market Microstructure XGBoost Pipeline\nThis module implements an XGBoost model with extensive market microstructure features\nand optional Optuna hyperparameter optimization followed by Ridge ensemble stacking\n\"\"\"\n\n# ===========================================================================================\n# PACKAGE INSTALLATION FOR MARKET MICROSTRUCTURE PIPELINE\n# ===========================================================================================\n\nimport subprocess\nimport sys\nimport os\nimport gc\nimport warnings\nwarnings.filterwarnings('ignore')\n\n# Install required packages for this pipeline\nprint(\"Installing packages for Market Microstructure XGBoost pipeline...\")\npackages_to_install = [\n    'xgboost==2.0.3',      # Specific version for stability\n    'optuna==3.5.0',       # Hyperparameter optimization\n    'pandas',\n    'numpy',\n    'scipy',\n    'scikit-learn',\n    'matplotlib'\n]\n\nfor package in packages_to_install:\n    subprocess.check_call([sys.executable, \"-m\", \"pip\", \"install\", package, \"--quiet\"])\n\nimport pandas as pd\nimport numpy as np\nfrom scipy.stats import pearsonr\nfrom typing import List, Dict, Tuple, Optional, Any\nfrom pathlib import Path\nimport json\nfrom sklearn.model_selection import KFold\nfrom sklearn.linear_model import Ridge\nimport matplotlib.pyplot as plt\n\n# Import packages after installation\nimport xgboost as xgb\nfrom xgboost import XGBRegressor\nimport optuna\n\nclass MarketMicrostructureConfiguration:\n    \"\"\"Configuration for Market Microstructure XGBoost Pipeline with model-specific settings\"\"\"\n    def __init__(self):\n        # Model identification\n        self.model_name = \"market_microstructure\"\n        self.model_directory = os.path.join(global_config.base_dir, \"market_microstructure_xgboost\")\n        \n        # Register model with global configuration\n        global_config.register_model(self.model_name, self.model_directory)\n        \n        # Data paths from global configuration\n        self.train_path = global_config.train_path\n        self.test_path = global_config.test_path\n        self.sample_sub_path = global_config.sample_sub_path\n        \n        # Model parameters\n        self.target = \"label\"\n        self.n_folds = 5\n        self.seed = 42\n        \n        # Optuna optimization settings\n        self.run_optuna = True\n        self.n_optuna_trials = 250\n        \n        # XGBoost parameters (optimized for market microstructure features)\n        self.xgb_params = {\n            \"tree_method\": \"hist\",  # Changed from gpu_hist for better compatibility\n            \"device\": \"cpu\",\n            \"colsample_bylevel\": 0.7,\n            \"colsample_bynode\": 0.7,\n            \"colsample_bytree\": 0.7,\n            \"gamma\": 1.5,\n            \"learning_rate\": 0.02,\n            \"max_depth\": 15,\n            \"max_leaves\": 20,\n            \"min_child_weight\": 10,\n            \"n_estimators\": 1500,\n            \"n_jobs\": -1,\n            \"random_state\": 42,\n            \"reg_alpha\": 30,\n            \"reg_lambda\": 60,\n            \"subsample\": 0.08,\n            \"verbosity\": 0\n        }\n        \n        # Output paths\n        self.intermediate_dir = os.path.join(self.model_directory, \"intermediate_predictions\")\n        self.submission_file = os.path.join(self.model_directory, \"submission.csv\")\n        self.metrics_file = os.path.join(self.model_directory, \"model_metrics.csv\")\n        self.feature_importance_file = os.path.join(self.model_directory, \"feature_importance.csv\")\n        self.optuna_results_file = os.path.join(self.model_directory, \"optuna_results.json\")\n        \n        # Model-specific columns to drop (these are typically redundant or problematic features)\n        self.cols_to_drop = [\n            'X697', 'X698', 'X699', 'X700', 'X701', 'X702', 'X703', 'X704', 'X705', 'X706', \n            'X707', 'X708', 'X709', 'X710', 'X711', 'X712', 'X713', 'X714', 'X715', 'X716',\n            'X717', 'X864', 'X867', 'X869', 'X870', 'X871', 'X872', 'X104', 'X110', 'X116',\n            'X122', 'X128', 'X134', 'X140', 'X146', 'X152', 'X158', 'X164', 'X170', 'X176',\n            'X182', 'X351', 'X357', 'X363', 'X369', 'X375', 'X381', 'X387', 'X393', 'X399',\n            'X405', 'X411', 'X417', 'X423', 'X429'\n        ]\n        \n        # Ensure directories exist\n        Path(self.model_directory).mkdir(parents=True, exist_ok=True)\n        Path(self.intermediate_dir).mkdir(parents=True, exist_ok=True)\n\nclass MarketMicrostructureDataProcessor:\n    \"\"\"Data processing utilities for Market Microstructure pipeline\"\"\"\n    @staticmethod\n    def reduce_mem_usage(dataframe: pd.DataFrame, dataset_name: str) -> pd.DataFrame:\n        \"\"\"Reduce memory usage by downcasting numeric types\"\"\"\n        print(f'Reducing memory usage for: {dataset_name}')\n        initial_mem_usage = dataframe.memory_usage().sum() / 1024**2\n        \n        for col in dataframe.columns:\n            col_type = dataframe[col].dtype\n            \n            if col_type != object:\n                c_min = dataframe[col].min()\n                c_max = dataframe[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                        dataframe[col] = dataframe[col].astype(np.int8)\n                    elif c_min > np.iinfo(np.int16).min and c_max < np.iinfo(np.int16).max:\n                        dataframe[col] = dataframe[col].astype(np.int16)\n                    elif c_min > np.iinfo(np.int32).min and c_max < np.iinfo(np.int32).max:\n                        dataframe[col] = dataframe[col].astype(np.int32)\n                    elif c_min > np.iinfo(np.int64).min and c_max < np.iinfo(np.int64).max:\n                        dataframe[col] = dataframe[col].astype(np.int64)\n                else:\n                    if c_min > np.finfo(np.float16).min and c_max < np.finfo(np.float16).max:\n                        dataframe[col] = dataframe[col].astype(np.float16)\n                    elif c_min > np.finfo(np.float32).min and c_max < np.finfo(np.float32).max:\n                        dataframe[col] = dataframe[col].astype(np.float32)\n                    else:\n                        dataframe[col] = dataframe[col].astype(np.float64)\n        \n        final_mem_usage = dataframe.memory_usage().sum() / 1024**2\n        print(f'Memory usage reduced by {100 * (initial_mem_usage - final_mem_usage) / initial_mem_usage:.1f}%')\n        \n        return dataframe\n\nclass MarketMicrostructureFeatureEngineer:\n    \"\"\"Feature engineering for market microstructure analysis\"\"\"\n    @staticmethod\n    def create_features(df: pd.DataFrame) -> pd.DataFrame:\n        \"\"\"Create extensive market microstructure features\"\"\"\n        df = df.copy()\n        \n        # Check if required columns exist\n        required_cols = ['bid_qty', 'ask_qty', 'buy_qty', 'sell_qty', 'volume']\n        available_cols = [col for col in required_cols if col in df.columns]\n        \n        if len(available_cols) < len(required_cols):\n            print(f\"Warning: Only {len(available_cols)} of {len(required_cols)} required columns available\")\n            print(f\"Missing columns: {set(required_cols) - set(available_cols)}\")\n            \n            # Create dummy columns if missing\n            for col in required_cols:\n                if col not in df.columns:\n                    df[col] = 0\n        \n        # Interaction features\n        df['bid_ask_interaction'] = df['bid_qty'] * df['ask_qty']\n        df['bid_buy_interaction'] = df['bid_qty'] * df['buy_qty']\n        df['bid_sell_interaction'] = df['bid_qty'] * df['sell_qty']\n        df['ask_buy_interaction'] = df['ask_qty'] * df['buy_qty']\n        df['ask_sell_interaction'] = df['ask_qty'] * df['sell_qty']\n        df['buy_sell_interaction'] = df['buy_qty'] * df['sell_qty']\n        \n        # Spread indicators\n        df['spread_indicator'] = (df['ask_qty'] - df['bid_qty']) / (df['ask_qty'] + df['bid_qty'] + 1e-8)\n        \n        # Volume-weighted features\n        df['volume_weighted_buy'] = df['buy_qty'] * df['volume']\n        df['volume_weighted_sell'] = df['sell_qty'] * df['volume']\n        df['volume_weighted_bid'] = df['bid_qty'] * df['volume']\n        df['volume_weighted_ask'] = df['ask_qty'] * df['volume']\n        \n        # Ratios\n        df['buy_sell_ratio'] = df['buy_qty'] / (df['sell_qty'] + 1e-8)\n        df['bid_ask_ratio'] = df['bid_qty'] / (df['ask_qty'] + 1e-8)\n        \n        # Order flow imbalance\n        df['order_flow_imbalance'] = (df['buy_qty'] - df['sell_qty']) / (df['volume'] + 1e-8)\n        \n        # Pressure indicators\n        df['buying_pressure'] = df['buy_qty'] / (df['volume'] + 1e-8)\n        df['selling_pressure'] = df['sell_qty'] / (df['volume'] + 1e-8)\n        \n        # Liquidity features\n        df['total_liquidity'] = df['bid_qty'] + df['ask_qty']\n        df['liquidity_imbalance'] = (df['bid_qty'] - df['ask_qty']) / (df['total_liquidity'] + 1e-8)\n        df['relative_spread'] = (df['ask_qty'] - df['bid_qty']) / (df['volume'] + 1e-8)\n        \n        # Trade intensity\n        df['trade_intensity'] = (df['buy_qty'] + df['sell_qty']) / (df['volume'] + 1e-8)\n        df['avg_trade_size'] = df['volume'] / (df['buy_qty'] + df['sell_qty'] + 1e-8)\n        df['net_trade_flow'] = (df['buy_qty'] - df['sell_qty']) / (df['buy_qty'] + df['sell_qty'] + 1e-8)\n        \n        # Market depth\n        df['depth_ratio'] = df['total_liquidity'] / (df['volume'] + 1e-8)\n        df['volume_participation'] = (df['buy_qty'] + df['sell_qty']) / (df['total_liquidity'] + 1e-8)\n        df['market_activity'] = df['volume'] * df['total_liquidity']\n        \n        # Spread proxy\n        df['effective_spread_proxy'] = np.abs(df['buy_qty'] - df['sell_qty']) / (df['volume'] + 1e-8)\n        df['realized_volatility_proxy'] = np.abs(df['order_flow_imbalance']) * df['volume']\n        \n        # Normalized volumes\n        df['normalized_buy_volume'] = df['buy_qty'] / (df['bid_qty'] + 1e-8)\n        df['normalized_sell_volume'] = df['sell_qty'] / (df['ask_qty'] + 1e-8)\n        \n        # Advanced features\n        df['liquidity_adjusted_imbalance'] = df['order_flow_imbalance'] * df['depth_ratio']\n        df['pressure_spread_interaction'] = df['buying_pressure'] * df['spread_indicator']\n        \n        # Additional market microstructure indicators\n        df['bid_ask_spread'] = df['ask_qty'] - df['bid_qty']\n        df['mid_price_proxy'] = (df['bid_qty'] + df['ask_qty']) / 2\n        df['price_pressure'] = df['net_trade_flow'] * df['volume']\n        df['liquidity_consumption'] = (df['buy_qty'] + df['sell_qty']) / (df['bid_qty'] + df['ask_qty'] + 1e-8)\n        \n        # Volatility and risk proxies\n        df['volume_volatility'] = df['volume'] * df['spread_indicator']\n        df['liquidity_risk'] = 1 / (df['total_liquidity'] + 1)\n        df['execution_risk'] = df['spread_indicator'] * df['liquidity_risk']\n        \n        # Clean up infinities and NaNs\n        df = df.replace([np.inf, -np.inf], np.nan)\n        df = df.fillna(0)\n        \n        print(f\"Created {len(df.columns)} features including engineered features\")\n        \n        return df\n\nclass MarketMicrostructurePipeline:\n    \"\"\"Complete Market Microstructure XGBoost Pipeline\"\"\"\n    def __init__(self, config: Optional[MarketMicrostructureConfiguration] = None):\n        self.config = config or MarketMicrostructureConfiguration()\n        self.data_processor = MarketMicrostructureDataProcessor()\n        self.feature_engineer = MarketMicrostructureFeatureEngineer()\n        \n    def optimize_ridge_hyperparameters(self, X_ensemble: pd.DataFrame, y: pd.Series) -> Dict[str, Any]:\n        \"\"\"Optimize Ridge regression hyperparameters using Optuna\"\"\"\n        print(\"Optimizing Ridge hyperparameters with Optuna...\")\n        \n        def objective(trial):\n            params = {\n                \"random_state\": self.config.seed,\n                \"alpha\": trial.suggest_float(\"alpha\", 0.001, 100),\n                \"tol\": trial.suggest_float(\"tol\", 1e-6, 1e-2),\n                \"solver\": trial.suggest_categorical(\"solver\", [\"auto\", \"svd\", \"cholesky\", \"lsqr\"])\n            }\n            \n            scores = []\n            kf = KFold(n_splits=self.config.n_folds, shuffle=True, random_state=self.config.seed)\n            \n            for train_idx, val_idx in kf.split(X_ensemble):\n                X_train, X_val = X_ensemble.iloc[train_idx], X_ensemble.iloc[val_idx]\n                y_train, y_val = y.iloc[train_idx], y.iloc[val_idx]\n                \n                model = Ridge(**params)\n                model.fit(X_train, y_train)\n                y_pred = model.predict(X_val)\n                \n                score = pearsonr(y_val, y_pred)[0]\n                scores.append(score)\n            \n            return np.mean(scores)\n        \n        # Set Optuna logging level\n        optuna.logging.set_verbosity(optuna.logging.WARNING)\n        \n        sampler = optuna.samplers.TPESampler(seed=self.config.seed, multivariate=True)\n        study = optuna.create_study(direction=\"maximize\", sampler=sampler)\n        study.optimize(objective, n_trials=self.config.n_optuna_trials, n_jobs=-1, catch=(ValueError,))\n        \n        best_params = study.best_params\n        print(f\"Best Ridge parameters: {best_params}\")\n        print(f\"Best cross-validation score: {study.best_value:.6f}\")\n        \n        ridge_params = {\n            \"random_state\": self.config.seed,\n            \"alpha\": best_params[\"alpha\"],\n            \"tol\": best_params[\"tol\"],\n            \"solver\": best_params.get(\"solver\", \"auto\")\n        }\n        \n        # Save Optuna results\n        optuna_results = {\n            \"best_params\": best_params,\n            \"best_value\": study.best_value,\n            \"n_trials\": len(study.trials),\n            \"optimization_history\": [trial.value for trial in study.trials if trial.value is not None]\n        }\n        \n        with open(self.config.optuna_results_file, 'w') as f:\n            json.dump(optuna_results, f, indent=2)\n        \n        # Register Optuna results with global configuration\n        global_config.register_model_output(\n            self.config.model_name,\n            'optuna_results',\n            self.config.optuna_results_file,\n            'analysis',\n            metadata=optuna_results\n        )\n        \n        return ridge_params\n    \n    def save_metrics(self, scores: Dict[str, float], ridge_params: Dict[str, Any]):\n        \"\"\"Save model metrics and parameters\"\"\"\n        metrics = {\n            \"model\": \"Market Microstructure XGBoost\",\n            \"xgboost_score\": scores.get(\"XGBoost\", 0),\n            \"ridge_ensemble_score\": scores.get(\"ridge_ensemble\", 0),\n            \"ridge_params\": ridge_params,\n            \"n_features\": scores.get(\"n_features\", 0),\n            \"n_folds\": self.config.n_folds\n        }\n        \n        metrics_df = pd.DataFrame([metrics])\n        metrics_df.to_csv(self.config.metrics_file, index=False)\n        print(f\"Metrics saved to {self.config.metrics_file}\")\n        \n        # Register metrics with global configuration\n        global_config.register_model_output(\n            self.config.model_name,\n            'model_metrics',\n            self.config.metrics_file,\n            'analysis',\n            metadata=metrics\n        )\n    \n    def run(self) -> Tuple[np.ndarray, Dict[str, float]]:\n        print(\"\\nStarting Market Microstructure XGBoost Pipeline...\")\n        print(\"This model focuses on market microstructure features and order flow dynamics\")\n        \n        # Update model status\n        global_config.update_model_status(self.config.model_name, \"running\")\n        \n        # Load data\n        print(\"\\nLoading data...\")\n        train = pd.read_parquet(self.config.train_path).reset_index(drop=True)\n        test = pd.read_parquet(self.config.test_path).reset_index(drop=True)\n        \n        print(f\"Original data shapes - Train: {train.shape}, Test: {test.shape}\")\n        \n        # Drop unnecessary columns\n        cols_to_drop_train = [col for col in self.config.cols_to_drop if col in train.columns]\n        cols_to_drop_test = [col for col in self.config.cols_to_drop + [\"label\"] if col in test.columns]\n        \n        if cols_to_drop_train:\n            train = train.drop(columns=cols_to_drop_train)\n            print(f\"Dropped {len(cols_to_drop_train)} columns from training data\")\n        \n        if cols_to_drop_test:\n            test = test.drop(columns=cols_to_drop_test)\n            print(f\"Dropped {len(cols_to_drop_test)} columns from test data\")\n        \n        # Reduce memory usage\n        train = self.data_processor.reduce_mem_usage(train, \"train\")\n        test = self.data_processor.reduce_mem_usage(test, \"test\")\n        \n        # Create market microstructure features\n        print(\"\\nEngineering market microstructure features...\")\n        train = self.feature_engineer.create_features(train)\n        test = self.feature_engineer.create_features(test)\n        \n        # Prepare data for modeling\n        X = train.drop(self.config.target, axis=1)\n        y = train[self.config.target]\n        X_test = test\n        \n        # Clean infinite values\n        X = X.replace([np.inf, -np.inf], np.nan).fillna(0)\n        X_test = X_test.replace([np.inf, -np.inf], np.nan).fillna(0)\n        \n        gc.collect()\n        \n        print(f\"\\nFinal data shapes - X: {X.shape}, X_test: {X_test.shape}\")\n        \n        # Initialize storage\n        scores = {\"n_features\": X.shape[1]}\n        oof_preds = {}\n        test_preds = {}\n        feature_importance_list = []\n        \n        print(\"\\nTraining XGBoost with Market Microstructure Features\")\n        print(f\"Using {self.config.n_folds}-fold cross-validation\")\n        \n        # Train XGBoost with cross-validation\n        oof = np.zeros(len(X))\n        test_predictions = np.zeros(len(X_test))\n        fold_scores = []\n        \n        kf = KFold(n_splits=self.config.n_folds, shuffle=True, random_state=self.config.seed)\n        \n        for fold, (train_idx, val_idx) in enumerate(kf.split(X)):\n            print(f\"\\nFold {fold + 1}/{self.config.n_folds}\")\n            \n            X_train, X_val = X.iloc[train_idx], X.iloc[val_idx]\n            y_train, y_val = y.iloc[train_idx], y.iloc[val_idx]\n            \n            print(f\"  Train samples: {len(X_train):,}, Val samples: {len(X_val):,}\")\n            \n            # Train model\n            model = XGBRegressor(**self.config.xgb_params)\n            model.fit(\n                X_train, y_train,\n                eval_set=[(X_val, y_val)],\n                early_stopping_rounds=50,\n                verbose=False\n            )\n            \n            # Make predictions\n            y_pred = model.predict(X_val)\n            oof[val_idx] = y_pred\n            \n            # Calculate fold score\n            score = pearsonr(y_val, y_pred)[0]\n            fold_scores.append(score)\n            print(f\"  Fold {fold + 1} Score: {score:.6f}\")\n            \n            # Test predictions\n            test_predictions += model.predict(X_test) / self.config.n_folds\n            \n            # Store feature importance\n            importance = model.feature_importances_\n            feature_importance_list.append(importance)\n            \n            # Clean up memory\n            del X_train, X_val, y_train, y_val, model\n            gc.collect()\n        \n        # Calculate overall score\n        overall_score = pearsonr(y, oof)[0]\n        \n        print(f\"\\nCross-validation completed:\")\n        print(f\"  Average fold score: {np.mean(fold_scores):.6f}\")\n        print(f\"  Overall OOF score: {overall_score:.6f}\")\n        \n        # Store results\n        oof_preds[\"XGBoost\"] = oof\n        test_preds[\"XGBoost\"] = test_predictions\n        scores[\"XGBoost\"] = overall_score\n        \n        # Ridge ensemble stacking\n        print(\"\\nCreating Ridge Regression Ensemble\")\n        \n        # Prepare ensemble data\n        X_ensemble = pd.DataFrame(oof_preds)\n        X_test_ensemble = pd.DataFrame(test_preds)\n        \n        # Optimize Ridge hyperparameters if enabled\n        if self.config.run_optuna:\n            ridge_params = self.optimize_ridge_hyperparameters(X_ensemble, y)\n        else:\n            ridge_params = {\"random_state\": self.config.seed, \"alpha\": 1.0}\n            print(\"Using default Ridge parameters (Optuna disabled)\")\n        \n        # Train Ridge ensemble with cross-validation\n        print(\"\\nTraining Ridge ensemble...\")\n        ridge_test_preds = np.zeros(len(X_test_ensemble))\n        ridge_oof_scores = []\n        \n        for fold, (train_idx, val_idx) in enumerate(kf.split(X_ensemble)):\n            X_train, X_val = X_ensemble.iloc[train_idx], X_ensemble.iloc[val_idx]\n            y_train, y_val = y.iloc[train_idx], y.iloc[val_idx]\n            \n            model = Ridge(**ridge_params)\n            model.fit(X_train, y_train)\n            \n            # Validation score\n            y_pred_val = model.predict(X_val)\n            fold_score = pearsonr(y_val, y_pred_val)[0]\n            ridge_oof_scores.append(fold_score)\n            \n            ridge_test_preds += model.predict(X_test_ensemble) / self.config.n_folds\n        \n        ridge_ensemble_score = np.mean(ridge_oof_scores)\n        print(f\"Ridge ensemble average validation score: {ridge_ensemble_score:.6f}\")\n        scores[\"ridge_ensemble\"] = ridge_ensemble_score\n        \n        # Save results\n        print(\"\\nSaving results...\")\n        \n        # Save submission\n        sub = pd.read_csv(self.config.sample_sub_path)\n        sub[\"prediction\"] = ridge_test_preds\n        sub.to_csv(self.config.submission_file, index=False)\n        print(f\"Submission saved to {self.config.submission_file}\")\n        \n        # Register submission with global configuration\n        global_config.register_model_output(\n            self.config.model_name,\n            'submission',\n            self.config.submission_file,\n            'submission',\n            metadata={\n                'xgboost_score': float(overall_score),\n                'ridge_ensemble_score': float(ridge_ensemble_score),\n                'n_features': X.shape[1],\n                'method': 'ridge_stacked_xgboost'\n            }\n        )\n        \n        # Save metrics\n        self.save_metrics(scores, ridge_params)\n        \n        # Save feature importance\n        avg_importance = np.mean(feature_importance_list, axis=0)\n        feature_importance = pd.DataFrame({\n            'feature': X.columns,\n            'importance': avg_importance\n        }).sort_values('importance', ascending=False)\n        \n        feature_importance.to_csv(self.config.feature_importance_file, index=False)\n        \n        # Register feature importance with global configuration\n        global_config.register_model_output(\n            self.config.model_name,\n            'feature_importance',\n            self.config.feature_importance_file,\n            'analysis',\n            metadata={\n                'top_10_features': feature_importance.head(10)['feature'].tolist(),\n                'top_10_importances': feature_importance.head(10)['importance'].tolist()\n            }\n        )\n        \n        print(\"\\nMarket Microstructure pipeline completed successfully!\")\n        print(f\"Final model score: {overall_score:.6f}\")\n        \n        return ridge_test_preds, scores\n\n# ===========================================================================================\n# EXECUTION BLOCK\n# ===========================================================================================\n\nif __name__ == \"__main__\":\n    try:\n        print(\"\\n🏛️ Running Market Microstructure XGBoost Pipeline\")\n        print(\"-\"*80)\n        \n        # Clean memory before starting\n        aggressive_memory_cleanup()\n        \n        # Create and run pipeline\n        market_pipeline = MarketMicrostructurePipeline()\n        predictions, scores = market_pipeline.run()\n        \n        # Extract final score (using XGBoost OOF score as primary metric)\n        final_score = scores.get(\"XGBoost\", 0.0)\n        \n        # Update global configuration with success\n        # FIXED: Use market_pipeline.config.model_name instead of self.config.model_name\n        global_config.update_model_status(market_pipeline.config.model_name, 'completed', score=final_score)\n        print(\"\\n✅ Market Microstructure pipeline completed successfully\")\n        \n    except Exception as e:\n        # Update global configuration with failure\n        error_msg = str(e)\n        global_config.update_model_status('market_microstructure', 'failed', error_message=error_msg)\n        print(f\"\\n❌ Market Microstructure pipeline failed: {error_msg}\")\n        raise\n        \n    finally:\n        # Clean up memory\n        aggressive_memory_cleanup()\n        \n        # Display current execution status\n        print(\"\\n\" + \"=\"*80)\n        print(\"Current Execution Status:\")\n        print(global_config.get_execution_summary())\n        print(\"=\"*80)","metadata":{"trusted":true,"jupyter":{"source_hidden":true}},"outputs":[],"execution_count":null},{"cell_type":"code","source":"# LIGHTGBM VOTING ENSEMBLE PIPELINE IMPLEMENTATION\n\n#!/usr/bin/env python\n# -*- coding: utf-8 -*-\n\"\"\"\nDRW Crypto Market Prediction - LightGBM Voting Ensemble Pipeline\nThis module implements a voting ensemble of LightGBM models trained on different\ntime-based folds with exponential time decay weighting for cryptocurrency prediction\n\"\"\"\n\n# ===========================================================================================\n# PACKAGE INSTALLATION FOR LIGHTGBM VOTING PIPELINE\n# ===========================================================================================\n\nimport subprocess\nimport sys\nimport os\nimport gc\nimport warnings\nwarnings.filterwarnings('ignore')\n\n# Install required packages for this pipeline\nprint(\"Installing packages for LightGBM Voting Ensemble pipeline...\")\npackages_to_install = [\n    'lightgbm==4.1.0',  # Specific version for stability\n    'pandas',\n    'numpy',\n    'scipy',\n    'scikit-learn',\n    'matplotlib'\n]\n\nfor package in packages_to_install:\n    subprocess.check_call([sys.executable, \"-m\", \"pip\", \"install\", package, \"--quiet\"])\n\n# ===========================================================================================\n# ===========================================================================================\n\nimport pandas as pd\nimport numpy as np\nfrom scipy.stats import pearsonr\nfrom typing import List, Dict, Tuple, Optional, Any\nfrom pathlib import Path\nimport json\nfrom sklearn.base import BaseEstimator, RegressorMixin\nimport matplotlib.pyplot as plt\n\n# Import LightGBM after installation\nimport lightgbm as lgb\n\nclass VotingModel(BaseEstimator, RegressorMixin):\n    \"\"\"Voting ensemble model that averages predictions from multiple estimators\"\"\"\n    def __init__(self, estimators):\n        super().__init__()\n        self.estimators = estimators\n\n    def fit(self, X, y=None):\n        return self\n\n    def predict(self, X):\n        y_preds = [estimator.predict(X) for estimator in self.estimators]\n        return np.mean(y_preds, axis=0)\n\n    def predict_proba(self, X):\n        y_preds = [estimator.predict_proba(X) for estimator in self.estimators]\n        return np.mean(y_preds, axis=0)\n\nclass LightGBMVotingConfiguration:\n    \"\"\"Configuration for LightGBM Voting Pipeline with model-specific settings\"\"\"\n    def __init__(self):\n        # Model identification\n        self.model_name = \"lightgbm_voting\"\n        self.model_directory = os.path.join(global_config.base_dir, \"lightgbm_voting\")\n        \n        # Register model with global configuration\n        global_config.register_model(self.model_name, self.model_directory)\n        \n        # Data paths from global configuration\n        self.train_path = global_config.train_path\n        self.test_path = global_config.test_path\n        self.sample_sub_path = global_config.sample_sub_path\n        \n        # Model-specific feature selection\n        self.feature_names = [\n            'X863', 'X856', 'X344', 'X598', 'X862', 'X385', 'X852', 'X603', 'X860', 'X674',\n            'X415', 'X345', 'X137', 'X855', 'X174', 'X302', 'X178', 'X532', 'X168', 'X612',\n            'bid_qty', 'ask_qty', 'buy_qty', 'sell_qty', 'volume',\n            'bid_ask_interaction', 'bid_buy_interaction', 'bid_sell_interaction',\n            'ask_buy_interaction', 'ask_sell_interaction'\n        ]\n        \n        # LightGBM parameters (optimized for financial time series)\n        self.lgb_params = {\n            \"boosting_type\": \"gbdt\",\n            \"objective\": \"regression\",\n            \"metric\": \"mae\",\n            \"colsample_bytree\": 0.55,\n            \"learning_rate\": 0.021,\n            \"min_child_samples\": 32,\n            \"min_child_weight\": 0.15,\n            \"max_depth\": -1,\n            \"n_jobs\": -1,\n            \"num_leaves\": 64,\n            \"random_state\": 42,\n            \"reg_alpha\": 80,\n            \"reg_lambda\": 100,\n            \"subsample\": 0.85,\n            \"verbosity\": 1,\n            \"device\": \"cpu\"  # Changed from \"gpu\" for better compatibility\n        }\n        \n        # Training parameters\n        self.decay_factor = 0.95\n        self.num_boost_round = 150\n        self.n_folds = 5\n        \n        # Time-based fold date ranges\n        self.fold_dates = {\n            1: ('2023-03-01 00:00:00', '2023-05-01 00:00:00'),\n            2: ('2023-05-01 00:00:00', '2023-07-01 00:00:00'),\n            3: ('2023-07-01 00:00:00', '2023-09-01 00:00:00'),\n            4: ('2023-09-01 00:00:00', '2023-11-01 00:00:00'),\n            5: ('2023-11-01 00:00:00', '2024-01-01 00:00:00'),\n            6: ('2024-01-01 00:00:00', '2024-03-01 00:00:00')\n        }\n        \n        # Output paths\n        self.intermediate_dir = os.path.join(self.model_directory, \"fold_models\")\n        self.submission_file = os.path.join(self.model_directory, \"submission.csv\")\n        self.metrics_file = os.path.join(self.model_directory, \"model_metrics.csv\")\n        self.feature_importance_file = os.path.join(self.model_directory, \"feature_importance.csv\")\n        self.fold_performance_file = os.path.join(self.model_directory, \"fold_performance.json\")\n        \n        # Ensure directories exist\n        Path(self.model_directory).mkdir(parents=True, exist_ok=True)\n        Path(self.intermediate_dir).mkdir(parents=True, exist_ok=True)\n\nclass LightGBMDataProcessor:\n    \"\"\"Data processing utilities for LightGBM pipeline\"\"\"\n    @staticmethod\n    def reduce_mem_usage(df: pd.DataFrame) -> pd.DataFrame:\n        \"\"\"Reduce memory usage by downcasting numeric types\"\"\"\n        start_mem = df.memory_usage().sum() / 1024**2\n        print(\"Memory usage of dataframe is {:.2f} MB\".format(start_mem))\n        \n        for col in df.columns:\n            col_type = df[col].dtype\n            if str(col_type) == \"category\":\n                continue\n                \n            if col_type != object:\n                c_min = df[col].min()\n                c_max = df[col].max()\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                    elif c_min > np.iinfo(np.int64).min and c_max < np.iinfo(np.int64).max:\n                        df[col] = df[col].astype(np.int64)\n                else:\n                    if c_min > np.finfo(np.float16).min and c_max < np.finfo(np.float16).max:\n                        df[col] = df[col].astype(np.float16)\n                    elif c_min > np.finfo(np.float32).min and c_max < np.finfo(np.float32).max:\n                        df[col] = df[col].astype(np.float32)\n                    else:\n                        df[col] = df[col].astype(np.float64)\n                        \n        end_mem = df.memory_usage().sum() / 1024**2\n        print(\"Memory usage after optimization is: {:.2f} MB\".format(end_mem))\n        print(\"Decreased by {:.1f}%\".format(100 * (start_mem - end_mem) / start_mem))\n        \n        return df\n    \n    @staticmethod\n    def create_time_weights(n_samples: int, decay_factor: float = 0.95) -> np.ndarray:\n        \"\"\"Create exponential decay weights giving more importance to recent data\"\"\"\n        positions = np.arange(n_samples)\n        normalized_positions = positions / (n_samples - 1)\n        weights = decay_factor ** (1 - normalized_positions)\n        weights = weights * n_samples / weights.sum()\n        return weights\n\nclass LightGBMFeatureEngineer:\n    \"\"\"Feature engineering for LightGBM pipeline\"\"\"\n    @staticmethod\n    def create_features(df: pd.DataFrame) -> pd.DataFrame:\n        \"\"\"Create market microstructure features\"\"\"\n        # Check if required columns exist\n        required_cols = ['bid_qty', 'ask_qty', 'buy_qty', 'sell_qty', 'volume']\n        available_cols = [col for col in required_cols if col in df.columns]\n        \n        if len(available_cols) < len(required_cols):\n            print(f\"Warning: Only {len(available_cols)} of {len(required_cols)} required columns available\")\n            # Create dummy columns if missing\n            for col in required_cols:\n                if col not in df.columns:\n                    df[col] = 0\n        \n        # Interaction features\n        df['bid_ask_interaction'] = df['bid_qty'] * df['ask_qty']\n        df['bid_buy_interaction'] = df['bid_qty'] * df['buy_qty']\n        df['bid_sell_interaction'] = df['bid_qty'] * df['sell_qty']\n        df['ask_buy_interaction'] = df['ask_qty'] * df['buy_qty']\n        df['ask_sell_interaction'] = df['ask_qty'] * df['sell_qty']\n        df['buy_sell_interaction'] = df['buy_qty'] * df['sell_qty']\n        \n        # Spread indicators\n        df['spread_indicator'] = (df['ask_qty'] - df['bid_qty']) / (df['ask_qty'] + df['bid_qty'] + 1e-8)\n        \n        # Volume-weighted features\n        df['volume_weighted_buy'] = df['buy_qty'] * df['volume']\n        df['volume_weighted_sell'] = df['sell_qty'] * df['volume']\n        df['volume_weighted_bid'] = df['bid_qty'] * df['volume']\n        df['volume_weighted_ask'] = df['ask_qty'] * df['volume']\n        \n        # Ratios\n        df['buy_sell_ratio'] = df['buy_qty'] / (df['sell_qty'] + 1e-8)\n        df['bid_ask_ratio'] = df['bid_qty'] / (df['ask_qty'] + 1e-8)\n        \n        # Order flow imbalance\n        df['order_flow_imbalance'] = (df['buy_qty'] - df['sell_qty']) / (df['volume'] + 1e-8)\n        \n        # Pressure indicators\n        df['buying_pressure'] = df['buy_qty'] / (df['volume'] + 1e-8)\n        df['selling_pressure'] = df['sell_qty'] / (df['volume'] + 1e-8)\n        \n        # Liquidity features\n        df['total_liquidity'] = df['bid_qty'] + df['ask_qty']\n        df['liquidity_imbalance'] = (df['bid_qty'] - df['ask_qty']) / (df['total_liquidity'] + 1e-8)\n        df['relative_spread'] = (df['ask_qty'] - df['bid_qty']) / (df['volume'] + 1e-8)\n        \n        # Trade intensity\n        df['trade_intensity'] = (df['buy_qty'] + df['sell_qty']) / (df['volume'] + 1e-8)\n        df['avg_trade_size'] = df['volume'] / (df['buy_qty'] + df['sell_qty'] + 1e-8)\n        df['net_trade_flow'] = (df['buy_qty'] - df['sell_qty']) / (df['buy_qty'] + df['sell_qty'] + 1e-8)\n        \n        # Market depth\n        df['depth_ratio'] = df['total_liquidity'] / (df['volume'] + 1e-8)\n        df['volume_participation'] = (df['buy_qty'] + df['sell_qty']) / (df['total_liquidity'] + 1e-8)\n        df['market_activity'] = df['volume'] * df['total_liquidity']\n        \n        # Spread proxy\n        df['effective_spread_proxy'] = np.abs(df['buy_qty'] - df['sell_qty']) / (df['volume'] + 1e-8)\n        df['realized_volatility_proxy'] = np.abs(df['order_flow_imbalance']) * df['volume']\n        \n        # Normalized volumes\n        df['normalized_buy_volume'] = df['buy_qty'] / (df['bid_qty'] + 1e-8)\n        df['normalized_sell_volume'] = df['sell_qty'] / (df['ask_qty'] + 1e-8)\n        \n        # Advanced features\n        df['liquidity_adjusted_imbalance'] = df['order_flow_imbalance'] * df['depth_ratio']\n        df['pressure_spread_interaction'] = df['buying_pressure'] * df['spread_indicator']\n        \n        # Clean up infinities and NaNs\n        df = df.replace([np.inf, -np.inf], np.nan)\n        df = df.fillna(0)\n        \n        return df\n\nclass LightGBMVotingPipeline:\n    \"\"\"Complete LightGBM Voting Ensemble Pipeline\"\"\"\n    def __init__(self, config: Optional[LightGBMVotingConfiguration] = None):\n        self.config = config or LightGBMVotingConfiguration()\n        self.data_processor = LightGBMDataProcessor()\n        self.feature_engineer = LightGBMFeatureEngineer()\n        \n    def assign_time_folds(self, train: pd.DataFrame) -> pd.DataFrame:\n        \"\"\"Assign time-based folds to training data\"\"\"\n        train['Fold'] = 0\n        \n        for fold_num, (start_date, end_date) in self.config.fold_dates.items():\n            mask = (train.index >= start_date) & (train.index < end_date)\n            train.loc[mask, 'Fold'] = fold_num\n            fold_size = mask.sum()\n            print(f\"Fold {fold_num}: {start_date} to {end_date} - {fold_size:,} samples\")\n        \n        return train\n    \n    def pearsonr_coeff(self, preds, data):\n        \"\"\"Custom evaluation metric for LightGBM\"\"\"\n        y_true = data.get_label()\n        valid_score = pearsonr(y_true, preds)[0]\n        return 'pearsonr_coeff_score', valid_score, True\n    \n    def train_single_model(self, train_data: lgb.Dataset, valid_data: lgb.Dataset) -> Tuple[lgb.Booster, float]:\n        \"\"\"Train a single LightGBM model\"\"\"\n        print(\"Training Model...\")\n        \n        model = lgb.train(\n            self.config.lgb_params,\n            train_data,\n            num_boost_round=self.config.num_boost_round,\n            valid_sets=[valid_data],\n            feval=self.pearsonr_coeff,\n            callbacks=[lgb.callback.log_evaluation(period=50)]\n        )\n        \n        valid_pred = model.predict(valid_data.get_data())\n        valid_score = pearsonr(valid_data.get_label(), valid_pred)[0]\n        print(f\"Validation Score: {valid_score:.6f}\")\n        \n        return model, valid_score\n    \n    def save_metrics(self, valid_scores: List[float], final_score: float):\n        \"\"\"Save pipeline metrics\"\"\"\n        metrics = {\n            \"model\": \"LightGBM Voting Ensemble\",\n            \"n_models\": len(valid_scores),\n            \"fold_scores\": valid_scores,\n            \"average_fold_score\": float(np.mean(valid_scores)),\n            \"std_fold_score\": float(np.std(valid_scores)),\n            \"final_ensemble_score\": float(final_score)\n        }\n        \n        metrics_df = pd.DataFrame([metrics])\n        metrics_df.to_csv(self.config.metrics_file, index=False)\n        print(f\"Metrics saved to {self.config.metrics_file}\")\n        \n        # Register metrics with global configuration\n        global_config.register_model_output(\n            self.config.model_name,\n            'model_metrics',\n            self.config.metrics_file,\n            'analysis',\n            metadata=metrics\n        )\n    \n    def save_fold_performance(self, fold_details: List[Dict[str, Any]]):\n        \"\"\"Save detailed fold performance data\"\"\"\n        with open(self.config.fold_performance_file, 'w') as f:\n            json.dump(fold_details, f, indent=2)\n        \n        # Register fold performance with global configuration\n        global_config.register_model_output(\n            self.config.model_name,\n            'fold_performance',\n            self.config.fold_performance_file,\n            'analysis',\n            metadata={'n_folds': len(fold_details)}\n        )\n    \n    def run(self) -> Tuple[np.ndarray, List[float]]:\n        print(\"\\nStarting LightGBM Voting Ensemble Pipeline...\")\n        print(\"This model uses time-based cross-validation with exponential decay weighting\")\n        \n        # Update model status\n        global_config.update_model_status(self.config.model_name, \"running\")\n        \n        # Load data\n        print(\"\\nLoading data...\")\n        train = pd.read_parquet(self.config.train_path)\n        test = pd.read_parquet(self.config.test_path)\n        \n        print(f\"Original data shapes - Train: {train.shape}, Test: {test.shape}\")\n        \n        # Reduce memory usage\n        train = self.data_processor.reduce_mem_usage(train)\n        test = self.data_processor.reduce_mem_usage(test)\n        \n        # Feature engineering\n        print(\"\\nEngineering features...\")\n        train = self.feature_engineer.create_features(train)\n        test = self.feature_engineer.create_features(test)\n        \n        # Assign time-based folds\n        print(\"\\nAssigning time-based folds...\")\n        train = self.assign_time_folds(train)\n        \n        # Create time weights\n        print(\"\\nCreating time decay weights...\")\n        train['weight'] = self.data_processor.create_time_weights(len(train), self.config.decay_factor)\n        \n        # Verify feature availability\n        available_features = [f for f in self.config.feature_names if f in train.columns]\n        if len(available_features) < len(self.config.feature_names):\n            print(f\"Warning: Only {len(available_features)} of {len(self.config.feature_names)} features available\")\n            print(f\"Missing features: {set(self.config.feature_names) - set(available_features)}\")\n        \n        # Initialize storage\n        models = []\n        valid_scores = []\n        fold_details = []\n        \n        # Train on each fold\n        print(\"\\nTraining models on time-based folds...\")\n        for fold in range(1, 6):  # Train 5 models, validate on fold 6\n            print(f\"\\n{'='*50}\")\n            print(f\"Training Fold {fold}/5\")\n            print(f\"{'='*50}\")\n            \n            # Prepare training and validation data\n            train_mask = train['Fold'] != fold\n            valid_mask = train['Fold'] == 6  # Always validate on the most recent fold\n            \n            X_train = train[train_mask][available_features]\n            w_train = train[train_mask]['weight']\n            X_valid = train[valid_mask][available_features]\n            w_valid = train[valid_mask]['weight']\n            y_train = train[train_mask]['label']\n            y_valid = train[valid_mask]['label']\n            \n            print(f\"Train samples: {len(X_train):,}\")\n            print(f\"Valid samples: {len(X_valid):,}\")\n            print(f\"Train time range: {X_train.index.min()} to {X_train.index.max()}\")\n            print(f\"Valid time range: {X_valid.index.min()} to {X_valid.index.max()}\")\n            \n            # Create LightGBM datasets\n            train_data = lgb.Dataset(\n                X_train, \n                label=y_train, \n                weight=w_train, \n                free_raw_data=False\n            ).construct()\n            \n            valid_data = lgb.Dataset(\n                X_valid, \n                label=y_valid, \n                weight=w_valid, \n                reference=train_data, \n                free_raw_data=False\n            ).construct()\n            \n            # Train model\n            model, valid_score = self.train_single_model(train_data, valid_data)\n            \n            # Save model\n            model_path = os.path.join(self.config.intermediate_dir, f'fold_{fold}_model.txt')\n            model.save_model(model_path)\n            print(f\"Model saved to {model_path}\")\n            \n            models.append(model)\n            valid_scores.append(valid_score)\n            \n            # Store fold details\n            fold_details.append({\n                'fold': fold,\n                'train_samples': len(X_train),\n                'valid_samples': len(X_valid),\n                'train_start': str(X_train.index.min()),\n                'train_end': str(X_train.index.max()),\n                'valid_start': str(X_valid.index.min()),\n                'valid_end': str(X_valid.index.max()),\n                'validation_score': float(valid_score)\n            })\n        \n        # Save fold performance details\n        self.save_fold_performance(fold_details)\n        \n        # Summary statistics\n        print(f\"\\n{'='*50}\")\n        print(\"Training Summary\")\n        print(f\"{'='*50}\")\n        print(f\"Individual fold scores: {[f'{score:.6f}' for score in valid_scores]}\")\n        print(f\"Average validation score: {np.mean(valid_scores):.6f}\")\n        print(f\"Standard deviation: {np.std(valid_scores):.6f}\")\n        \n        # Create voting ensemble\n        print(\"\\nCreating voting ensemble...\")\n        lgbm_voting = VotingModel(models)\n        \n        # Make predictions\n        print(\"Making predictions on test data...\")\n        test_features = [f for f in available_features if f in test.columns]\n        predictions = lgbm_voting.predict(test[test_features])\n        \n        # Save submission\n        submission = pd.read_csv(self.config.sample_sub_path)\n        submission[\"prediction\"] = predictions\n        submission.to_csv(self.config.submission_file, index=False)\n        print(f\"Submission saved to {self.config.submission_file}\")\n        \n        # Register submission with global configuration\n        final_score = np.mean(valid_scores)\n        global_config.register_model_output(\n            self.config.model_name,\n            'submission',\n            self.config.submission_file,\n            'submission',\n            metadata={\n                'ensemble_score': float(final_score),\n                'n_models': len(models),\n                'fold_scores': [float(s) for s in valid_scores],\n                'method': 'time_based_voting_ensemble'\n            }\n        )\n        \n        # Save metrics\n        self.save_metrics(valid_scores, final_score)\n        \n        # Create feature importance summary\n        print(\"\\nCalculating feature importance...\")\n        feature_importance_sum = np.zeros(len(available_features))\n        for model in models:\n            feature_importance_sum += model.feature_importance(importance_type='gain')\n        \n        feature_importance_df = pd.DataFrame({\n            'feature': available_features,\n            'importance': feature_importance_sum / len(models)\n        }).sort_values('importance', ascending=False)\n        \n        feature_importance_df.to_csv(self.config.feature_importance_file, index=False)\n        print(f\"Feature importance saved to {self.config.feature_importance_file}\")\n        \n        # Register feature importance with global configuration\n        global_config.register_model_output(\n            self.config.model_name,\n            'feature_importance',\n            self.config.feature_importance_file,\n            'analysis',\n            metadata={\n                'top_10_features': feature_importance_df.head(10)['feature'].tolist(),\n                'top_10_importances': feature_importance_df.head(10)['importance'].tolist()\n            }\n        )\n        \n        print(f\"\\nTop 10 most important features:\")\n        for idx, row in feature_importance_df.head(10).iterrows():\n            print(f\"  {row['feature']}: {row['importance']:.2f}\")\n        \n        print(\"\\nLightGBM Voting pipeline completed successfully!\")\n        print(f\"Final ensemble score: {final_score:.6f}\")\n        \n        return predictions, valid_scores\n\n# ===========================================================================================\n# EXECUTION BLOCK\n# ===========================================================================================\n\nif __name__ == \"__main__\":\n    try:\n        print(\"\\n🌲 Running LightGBM Voting Ensemble Pipeline\")\n        print(\"-\"*80)\n        \n        # Clean memory before starting\n        aggressive_memory_cleanup()\n        \n        # Create and run pipeline\n        lgbm_pipeline = LightGBMVotingPipeline()\n        predictions, valid_scores = lgbm_pipeline.run()\n        \n        # Extract final score\n        final_score = np.mean(valid_scores)\n        \n        # Update global configuration with success\n        global_config.update_model_status('lightgbm_voting', 'completed', score=final_score)\n        print(\"\\n✅ LightGBM Voting pipeline completed successfully\")\n        \n    except Exception as e:\n        # Update global configuration with failure\n        error_msg = str(e)\n        global_config.update_model_status('lightgbm_voting', 'failed', error_message=error_msg)\n        print(f\"\\n❌ LightGBM Voting pipeline failed: {error_msg}\")\n        raise\n        \n    finally:\n        # Clean up memory\n        aggressive_memory_cleanup()\n        \n        # Display current execution status\n        print(\"\\n\" + \"=\"*80)\n        print(\"Current Execution Status:\")\n        print(global_config.get_execution_summary())\n        print(\"=\"*80)","metadata":{"trusted":true,"jupyter":{"source_hidden":true,"outputs_hidden":true},"execution":{"iopub.status.busy":"2025-06-04T16:09:38.792632Z","iopub.execute_input":"2025-06-04T16:09:38.792951Z","iopub.status.idle":"2025-06-04T16:12:00.642285Z","shell.execute_reply.started":"2025-06-04T16:09:38.792933Z","shell.execute_reply":"2025-06-04T16:12:00.639165Z"},"collapsed":true},"outputs":[],"execution_count":null},{"cell_type":"code","source":"# Final Ensemble Building with Meta-Learning\n# !/usr/bin/env python\n# -*- coding: utf-8 -*-\n\"\"\"\nDRW Crypto Market Prediction - Advanced Meta-Learning Ensemble Builder\nThis module creates the final ensemble using sophisticated meta-learning techniques\nincluding multi-level stacking, dynamic model selection, and adaptive weighting\n\"\"\"\n\nimport subprocess\nimport sys\nimport os\nimport gc\nimport warnings\nimport json\nimport pickle\nimport pandas as pd\nimport numpy as np\nfrom scipy.stats import pearsonr, spearmanr\nfrom typing import List, Dict, Tuple, Optional, Any, Union\nfrom pathlib import Path\nfrom sklearn.model_selection import KFold, TimeSeriesSplit\nfrom sklearn.preprocessing import StandardScaler, RobustScaler\nfrom sklearn.linear_model import Ridge, ElasticNet, HuberRegressor\nfrom sklearn.ensemble import RandomForestRegressor, ExtraTreesRegressor\nfrom sklearn.neural_network import MLPRegressor\nfrom sklearn.cluster import KMeans\nfrom sklearn.decomposition import PCA\nimport matplotlib.pyplot as plt\nimport seaborn as sns\nfrom abc import ABC, abstractmethod\n\nwarnings.filterwarnings('ignore')\n\n# Install required packages for meta-learning ensemble\nprint(\"Installing packages for Meta-Learning Ensemble Builder...\")\npackages_to_install = [\n    'flaml==2.1.1',\n    'lightgbm==4.1.0',\n    'optuna==3.4.0',\n    'scikit-optimize==0.9.0',\n    'matplotlib',\n    'seaborn'\n]\n\nfor package in packages_to_install:\n    try:\n        subprocess.check_call([sys.executable, \"-m\", \"pip\", \"install\", package, \"--quiet\"])\n    except:\n        print(f\"Warning: Could not install {package}\")\n\n# =============================================================================\n# Base Meta-Learning Classes\n# =============================================================================\n\nclass BaseMetaLearner(ABC):\n    \"\"\"Abstract base class for meta-learners\"\"\"\n    \n    @abstractmethod\n    def fit(self, predictions: Dict[str, np.ndarray], y: np.ndarray, **kwargs):\n        \"\"\"Fit the meta-learner\"\"\"\n        pass\n    \n    @abstractmethod\n    def predict(self, predictions: Dict[str, np.ndarray], **kwargs) -> np.ndarray:\n        \"\"\"Make predictions using the meta-learner\"\"\"\n        pass\n    \n    @abstractmethod\n    def get_feature_importance(self) -> Dict[str, float]:\n        \"\"\"Get feature importance scores\"\"\"\n        pass\n\n# =============================================================================\n# Multi-Level Stacking Meta-Learner\n# =============================================================================\n\nclass MultiLevelStackingMetaLearner(BaseMetaLearner):\n    \"\"\"Multi-level stacking with diverse meta-features and models\"\"\"\n    \n    def __init__(self, n_levels: int = 2, random_state: int = 42):\n        self.n_levels = n_levels\n        self.random_state = random_state\n        self.level_models = {}\n        self.feature_generators = {}\n        self.feature_names = []\n        self.scaler = RobustScaler()\n        \n    def _create_polynomial_features(self, predictions: Dict[str, np.ndarray]) -> Dict[str, np.ndarray]:\n        \"\"\"Create polynomial interaction features\"\"\"\n        poly_features = {}\n        model_names = list(predictions.keys())\n        \n        # Quadratic features\n        for i, name1 in enumerate(model_names):\n            poly_features[f'{name1}_squared'] = predictions[name1] ** 2\n            \n            for j, name2 in enumerate(model_names[i+1:], i+1):\n                poly_features[f'{name1}_x_{name2}'] = predictions[name1] * predictions[name2]\n        \n        return poly_features\n    \n    def _create_rank_features(self, predictions: Dict[str, np.ndarray]) -> Dict[str, np.ndarray]:\n        \"\"\"Create rank-based features\"\"\"\n        rank_features = {}\n        pred_array = np.column_stack(list(predictions.values()))\n        \n        # Rank of each prediction\n        for i, name in enumerate(predictions.keys()):\n            rank_features[f'{name}_rank'] = np.argsort(np.argsort(predictions[name])) / len(predictions[name])\n        \n        # Consensus ranking\n        mean_ranks = np.mean([rank_features[f'{name}_rank'] for name in predictions.keys()], axis=0)\n        rank_features['consensus_rank'] = mean_ranks\n        \n        return rank_features\n    \n    def _create_statistical_features(self, predictions: Dict[str, np.ndarray]) -> Dict[str, np.ndarray]:\n        \"\"\"Create statistical aggregation features\"\"\"\n        pred_array = np.column_stack(list(predictions.values()))\n        \n        stats_features = {\n            'mean': np.mean(pred_array, axis=1),\n            'median': np.median(pred_array, axis=1),\n            'std': np.std(pred_array, axis=1),\n            'min': np.min(pred_array, axis=1),\n            'max': np.max(pred_array, axis=1),\n            'range': np.ptp(pred_array, axis=1),\n            'iqr': np.percentile(pred_array, 75, axis=1) - np.percentile(pred_array, 25, axis=1),\n            'skew': np.array([np.mean((pred_array[i] - np.mean(pred_array[i])) ** 3) / \n                             (np.std(pred_array[i]) ** 3 + 1e-8) for i in range(len(pred_array))]),\n            'kurtosis': np.array([np.mean((pred_array[i] - np.mean(pred_array[i])) ** 4) / \n                                 (np.std(pred_array[i]) ** 4 + 1e-8) - 3 for i in range(len(pred_array))])\n        }\n        \n        # Coefficient of variation\n        stats_features['cv'] = np.where(stats_features['mean'] != 0,\n                                       stats_features['std'] / np.abs(stats_features['mean']),\n                                       0)\n        \n        return stats_features\n    \n    def _create_distance_features(self, predictions: Dict[str, np.ndarray]) -> Dict[str, np.ndarray]:\n        \"\"\"Create distance-based features\"\"\"\n        distance_features = {}\n        pred_array = np.column_stack(list(predictions.values()))\n        \n        # Distance from mean\n        mean_pred = np.mean(pred_array, axis=1)\n        for i, name in enumerate(predictions.keys()):\n            distance_features[f'{name}_dist_from_mean'] = np.abs(predictions[name] - mean_pred)\n        \n        # Pairwise distances\n        model_names = list(predictions.keys())\n        for i in range(len(model_names)):\n            for j in range(i + 1, len(model_names)):\n                distance_features[f'dist_{model_names[i]}_{model_names[j]}'] = \\\n                    np.abs(predictions[model_names[i]] - predictions[model_names[j]])\n        \n        return distance_features\n    \n    def _generate_meta_features(self, predictions: Dict[str, np.ndarray]) -> np.ndarray:\n        \"\"\"Generate comprehensive meta-features\"\"\"\n        all_features = {}\n        \n        # Base predictions\n        all_features.update(predictions)\n        \n        # Polynomial features\n        all_features.update(self._create_polynomial_features(predictions))\n        \n        # Rank features\n        all_features.update(self._create_rank_features(predictions))\n        \n        # Statistical features\n        all_features.update(self._create_statistical_features(predictions))\n        \n        # Distance features\n        all_features.update(self._create_distance_features(predictions))\n        \n        # Store feature names for later use\n        self.feature_names = list(all_features.keys())\n        \n        # Convert to array\n        return np.column_stack([all_features[name] for name in self.feature_names])\n    \n    def fit(self, predictions: Dict[str, np.ndarray], y: np.ndarray, **kwargs):\n        \"\"\"Fit multi-level stacking models\"\"\"\n        print(\"\\nTraining Multi-Level Stacking Meta-Learner...\")\n        \n        # Generate meta-features\n        X_meta = self._generate_meta_features(predictions)\n        X_meta_scaled = self.scaler.fit_transform(X_meta)\n        \n        # Level 1: Diverse base meta-learners\n        print(\"  Training Level 1 models...\")\n        \n        # Import LightGBM after installation\n        try:\n            import lightgbm as lgb\n            \n            self.level_models['level1'] = {\n                'rf': RandomForestRegressor(\n                    n_estimators=200,\n                    max_depth=10,\n                    min_samples_split=20,\n                    random_state=self.random_state\n                ),\n                'et': ExtraTreesRegressor(\n                    n_estimators=200,\n                    max_depth=10,\n                    min_samples_split=20,\n                    random_state=self.random_state\n                ),\n                'lgb': lgb.LGBMRegressor(\n                    n_estimators=200,\n                    max_depth=8,\n                    learning_rate=0.05,\n                    random_state=self.random_state,\n                    verbose=-1\n                ),\n                'mlp': MLPRegressor(\n                    hidden_layer_sizes=(100, 50),\n                    activation='relu',\n                    solver='adam',\n                    alpha=0.01,\n                    random_state=self.random_state,\n                    max_iter=500\n                ),\n                'ridge': Ridge(alpha=1.0, random_state=self.random_state),\n                'elastic': ElasticNet(alpha=0.1, l1_ratio=0.5, random_state=self.random_state),\n                'huber': HuberRegressor(epsilon=1.35, alpha=0.01)\n            }\n        except ImportError:\n            print(\"  Warning: LightGBM not available, using alternative models\")\n            self.level_models['level1'] = {\n                'rf': RandomForestRegressor(n_estimators=200, max_depth=10, random_state=self.random_state),\n                'et': ExtraTreesRegressor(n_estimators=200, max_depth=10, random_state=self.random_state),\n                'ridge': Ridge(alpha=1.0, random_state=self.random_state),\n                'elastic': ElasticNet(alpha=0.1, l1_ratio=0.5, random_state=self.random_state)\n            }\n        \n        # Train level 1 models\n        level1_predictions = {}\n        for name, model in self.level_models['level1'].items():\n            model.fit(X_meta_scaled, y)\n            level1_predictions[name] = model.predict(X_meta_scaled)\n            score = pearsonr(level1_predictions[name], y)[0]\n            print(f\"    {name}: correlation = {score:.4f}\")\n        \n        # Level 2: Meta-blender\n        if self.n_levels >= 2:\n            print(\"  Training Level 2 blender...\")\n            \n            # Combine level 1 predictions with original predictions\n            level2_features = np.column_stack(\n                list(level1_predictions.values()) + \n                [predictions[name] for name in predictions.keys()]\n            )\n            \n            # Train final blender\n            self.level_models['blender'] = Ridge(alpha=0.1, random_state=self.random_state)\n            self.level_models['blender'].fit(level2_features, y)\n            \n            final_pred = self.level_models['blender'].predict(level2_features)\n            final_score = pearsonr(final_pred, y)[0]\n            print(f\"    Final blender: correlation = {final_score:.4f}\")\n    \n    def predict(self, predictions: Dict[str, np.ndarray], **kwargs) -> np.ndarray:\n        \"\"\"Make predictions using stacked models\"\"\"\n        # Generate meta-features\n        X_meta = self._generate_meta_features(predictions)\n        X_meta_scaled = self.scaler.transform(X_meta)\n        \n        # Level 1 predictions\n        level1_predictions = {}\n        for name, model in self.level_models['level1'].items():\n            level1_predictions[name] = model.predict(X_meta_scaled)\n        \n        # Final prediction\n        if 'blender' in self.level_models:\n            level2_features = np.column_stack(\n                list(level1_predictions.values()) + \n                [predictions[name] for name in predictions.keys()]\n            )\n            return self.level_models['blender'].predict(level2_features)\n        else:\n            # Simple average if no blender\n            return np.mean(list(level1_predictions.values()), axis=0)\n    \n    def get_feature_importance(self) -> Dict[str, float]:\n        \"\"\"Get aggregated feature importance\"\"\"\n        importance_scores = {}\n        \n        # Get importance from tree-based models\n        for name, model in self.level_models['level1'].items():\n            if hasattr(model, 'feature_importances_'):\n                for i, feat_name in enumerate(self.feature_names):\n                    if feat_name not in importance_scores:\n                        importance_scores[feat_name] = 0\n                    importance_scores[feat_name] += model.feature_importances_[i]\n        \n        # Normalize\n        total = sum(importance_scores.values())\n        if total > 0:\n            importance_scores = {k: v/total for k, v in importance_scores.items()}\n        \n        return importance_scores\n\n# =============================================================================\n# Dynamic Model Selection Meta-Learner\n# =============================================================================\n\nclass DynamicModelSelector(BaseMetaLearner):\n    \"\"\"Selects models dynamically based on market regime detection\"\"\"\n    \n    def __init__(self, n_regimes: int = 4, lookback_window: int = 100, random_state: int = 42):\n        self.n_regimes = n_regimes\n        self.lookback_window = lookback_window\n        self.random_state = random_state\n        self.regime_detector = None\n        self.regime_models = {}\n        self.regime_features_scaler = StandardScaler()\n        self.model_performance_history = {}\n    \n    def _extract_regime_features(self, predictions: Dict[str, np.ndarray]) -> np.ndarray:\n        \"\"\"Extract features for regime detection\"\"\"\n        pred_array = np.column_stack(list(predictions.values()))\n        \n        regime_features = []\n        \n        # Model disagreement\n        model_std = np.std(pred_array, axis=1)\n        regime_features.append(model_std)\n        \n        # Prediction level\n        mean_pred = np.mean(pred_array, axis=1)\n        regime_features.append(mean_pred)\n        \n        # Rolling statistics (simulated with expanding windows)\n        window_size = min(self.lookback_window, len(mean_pred) // 10)\n        if window_size > 1:\n            # Rolling volatility proxy\n            rolling_std = np.array([\n                np.std(mean_pred[max(0, i-window_size):i+1]) \n                for i in range(len(mean_pred))\n            ])\n            regime_features.append(rolling_std)\n            \n            # Trend strength proxy\n            trend_strength = np.array([\n                (mean_pred[i] - mean_pred[max(0, i-window_size)]) / (window_size + 1e-8)\n                for i in range(len(mean_pred))\n            ])\n            regime_features.append(trend_strength)\n        \n        # Model correlation changes\n        for i in range(len(predictions)):\n            for j in range(i + 1, len(predictions)):\n                model_names = list(predictions.keys())\n                corr_proxy = predictions[model_names[i]] * predictions[model_names[j]]\n                regime_features.append(corr_proxy)\n        \n        return np.column_stack(regime_features)\n    \n    def _detect_regimes(self, regime_features: np.ndarray) -> np.ndarray:\n        \"\"\"Detect market regimes using clustering\"\"\"\n        # Apply PCA for dimensionality reduction\n        pca = PCA(n_components=min(5, regime_features.shape[1]), random_state=self.random_state)\n        regime_features_pca = pca.fit_transform(regime_features)\n        \n        # Cluster into regimes\n        self.regime_detector = KMeans(\n            n_clusters=self.n_regimes,\n            random_state=self.random_state,\n            n_init=10\n        )\n        regimes = self.regime_detector.fit_predict(regime_features_pca)\n        \n        return regimes\n    \n    def _select_best_models_for_regime(self, predictions: Dict[str, np.ndarray], \n                                     y: np.ndarray, regime_mask: np.ndarray) -> Dict[str, float]:\n        \"\"\"Select best performing models for a specific regime\"\"\"\n        regime_predictions = {name: pred[regime_mask] for name, pred in predictions.items()}\n        y_regime = y[regime_mask]\n        \n        model_scores = {}\n        for name, pred in regime_predictions.items():\n            if len(pred) > 10:  # Minimum samples\n                score = pearsonr(pred, y_regime)[0]\n                model_scores[name] = max(0, score)  # Clip negative correlations\n            else:\n                model_scores[name] = 0\n        \n        # Normalize scores\n        total_score = sum(model_scores.values())\n        if total_score > 0:\n            model_scores = {k: v/total_score for k, v in model_scores.items()}\n        else:\n            # Equal weights if all models perform poorly\n            model_scores = {k: 1/len(predictions) for k in predictions.keys()}\n        \n        return model_scores\n    \n    def fit(self, predictions: Dict[str, np.ndarray], y: np.ndarray, **kwargs):\n        \"\"\"Fit regime-specific models\"\"\"\n        print(\"\\nTraining Dynamic Model Selector...\")\n        \n        # Extract regime features\n        regime_features = self._extract_regime_features(predictions)\n        regime_features_scaled = self.regime_features_scaler.fit_transform(regime_features)\n        \n        # Detect regimes\n        regimes = self._detect_regimes(regime_features_scaled)\n        \n        print(f\"  Detected {self.n_regimes} market regimes\")\n        unique_regimes, regime_counts = np.unique(regimes, return_counts=True)\n        for regime, count in zip(unique_regimes, regime_counts):\n            print(f\"    Regime {regime}: {count} samples ({count/len(regimes)*100:.1f}%)\")\n        \n        # Train regime-specific models\n        for regime in unique_regimes:\n            regime_mask = regimes == regime\n            \n            if np.sum(regime_mask) > 20:  # Minimum samples per regime\n                # Select best models for this regime\n                model_weights = self._select_best_models_for_regime(predictions, y, regime_mask)\n                \n                # Create weighted ensemble for regime\n                self.regime_models[regime] = {\n                    'weights': model_weights,\n                    'n_samples': np.sum(regime_mask)\n                }\n                \n                # Display regime model selection\n                print(f\"  Regime {regime} model weights:\")\n                for model_name, weight in sorted(model_weights.items(), key=lambda x: x[1], reverse=True):\n                    if weight > 0.1:  # Only show significant weights\n                        print(f\"    {model_name}: {weight:.3f}\")\n    \n    def predict(self, predictions: Dict[str, np.ndarray], **kwargs) -> np.ndarray:\n        \"\"\"Predict using regime-appropriate models\"\"\"\n        # Extract regime features\n        regime_features = self._extract_regime_features(predictions)\n        regime_features_scaled = self.regime_features_scaler.transform(regime_features)\n        \n        # Detect current regimes\n        regime_features_pca = PCA(n_components=min(5, regime_features_scaled.shape[1]), \n                                 random_state=self.random_state).fit_transform(regime_features_scaled)\n        current_regimes = self.regime_detector.predict(regime_features_pca)\n        \n        # Make predictions based on regime\n        final_predictions = np.zeros(len(next(iter(predictions.values()))))\n        \n        for regime in np.unique(current_regimes):\n            regime_mask = current_regimes == regime\n            \n            if regime in self.regime_models:\n                weights = self.regime_models[regime]['weights']\n                regime_pred = np.zeros(np.sum(regime_mask))\n                \n                for model_name, weight in weights.items():\n                    regime_pred += weight * predictions[model_name][regime_mask]\n                \n                final_predictions[regime_mask] = regime_pred\n            else:\n                # Fallback to equal weights\n                regime_pred = np.mean([pred[regime_mask] for pred in predictions.values()], axis=0)\n                final_predictions[regime_mask] = regime_pred\n        \n        return final_predictions\n    \n    def get_feature_importance(self) -> Dict[str, float]:\n        \"\"\"Get model importance across regimes\"\"\"\n        importance = {}\n        \n        for regime, regime_info in self.regime_models.items():\n            weights = regime_info['weights']\n            n_samples = regime_info['n_samples']\n            \n            for model_name, weight in weights.items():\n                if model_name not in importance:\n                    importance[model_name] = 0\n                importance[model_name] += weight * n_samples\n        \n        # Normalize\n        total = sum(importance.values())\n        if total > 0:\n            importance = {k: v/total for k, v in importance.items()}\n        \n        return importance\n\n# =============================================================================\n# Bayesian Model Averaging Meta-Learner\n# =============================================================================\n\nclass BayesianModelAveraging(BaseMetaLearner):\n    \"\"\"Bayesian approach to model combination with uncertainty quantification\"\"\"\n    \n    def __init__(self, prior_strength: float = 1.0, mcmc_samples: int = 1000, random_state: int = 42):\n        self.prior_strength = prior_strength\n        self.mcmc_samples = mcmc_samples\n        self.random_state = random_state\n        self.posterior_weights = None\n        self.model_precisions = None\n        self.evidence = None\n    \n    def _compute_model_likelihood(self, predictions: np.ndarray, y: np.ndarray, \n                                 precision: float) -> float:\n        \"\"\"Compute likelihood of data given model predictions\"\"\"\n        residuals = y - predictions\n        log_likelihood = -0.5 * precision * np.sum(residuals ** 2)\n        log_likelihood += 0.5 * len(y) * (np.log(precision) - np.log(2 * np.pi))\n        return log_likelihood\n    \n    def _sample_posterior(self, predictions: Dict[str, np.ndarray], y: np.ndarray):\n        \"\"\"Sample from posterior distribution using simple MCMC\"\"\"\n        n_models = len(predictions)\n        model_names = list(predictions.keys())\n        pred_matrix = np.column_stack([predictions[name] for name in model_names])\n        \n        # Initialize\n        weights = np.ones(n_models) / n_models\n        precisions = np.ones(n_models)\n        \n        # Storage for samples\n        weight_samples = []\n        precision_samples = []\n        \n        # Simple Metropolis-Hastings\n        for iteration in range(self.mcmc_samples * 2):  # Extra samples for burn-in\n            # Update weights (Dirichlet-like proposal)\n            proposed_weights = np.random.dirichlet(weights * 10 + self.prior_strength)\n            \n            # Compute ensemble predictions\n            current_pred = np.dot(pred_matrix, weights)\n            proposed_pred = np.dot(pred_matrix, proposed_weights)\n            \n            # Compute acceptance ratio\n            current_likelihood = self._compute_model_likelihood(current_pred, y, np.mean(precisions))\n            proposed_likelihood = self._compute_model_likelihood(proposed_pred, y, np.mean(precisions))\n            \n            accept_prob = np.exp(proposed_likelihood - current_likelihood)\n            \n            if np.random.rand() < accept_prob:\n                weights = proposed_weights\n            \n            # Update precisions (Gamma-like)\n            for i in range(n_models):\n                residuals = y - pred_matrix[:, i]\n                shape = len(y) / 2 + 1\n                scale = 2 / (np.sum(residuals ** 2) + 1e-6)\n                precisions[i] = np.random.gamma(shape, scale)\n            \n            # Store samples after burn-in\n            if iteration >= self.mcmc_samples:\n                weight_samples.append(weights.copy())\n                precision_samples.append(precisions.copy())\n        \n        # Compute posterior statistics\n        self.posterior_weights = np.mean(weight_samples, axis=0)\n        self.model_precisions = np.mean(precision_samples, axis=0)\n        \n        # Store samples for uncertainty quantification\n        self.weight_samples = np.array(weight_samples)\n        self.precision_samples = np.array(precision_samples)\n    \n    def fit(self, predictions: Dict[str, np.ndarray], y: np.ndarray, **kwargs):\n        \"\"\"Fit Bayesian model averaging\"\"\"\n        print(\"\\nTraining Bayesian Model Averaging...\")\n        \n        # Sample from posterior\n        self._sample_posterior(predictions, y)\n        \n        # Display results\n        print(\"  Posterior model weights:\")\n        for name, weight in zip(predictions.keys(), self.posterior_weights):\n            print(f\"    {name}: {weight:.3f} (precision: {self.model_precisions[list(predictions.keys()).index(name)]:.2f})\")\n        \n        # Compute model evidence (marginal likelihood approximation)\n        pred_matrix = np.column_stack(list(predictions.values()))\n        ensemble_pred = np.dot(pred_matrix, self.posterior_weights)\n        self.evidence = pearsonr(ensemble_pred, y)[0]\n        print(f\"  Model evidence (correlation): {self.evidence:.4f}\")\n    \n    def predict(self, predictions: Dict[str, np.ndarray], **kwargs) -> np.ndarray:\n        \"\"\"Make predictions with uncertainty quantification\"\"\"\n        pred_matrix = np.column_stack([predictions[name] for name in predictions.keys()])\n        \n        # Point prediction\n        point_prediction = np.dot(pred_matrix, self.posterior_weights)\n        \n        # Uncertainty quantification if requested\n        if kwargs.get('return_uncertainty', False):\n            # Sample predictions\n            sample_predictions = []\n            n_uncertainty_samples = min(100, len(self.weight_samples))\n            \n            for i in range(n_uncertainty_samples):\n                sample_pred = np.dot(pred_matrix, self.weight_samples[i])\n                sample_predictions.append(sample_pred)\n            \n            sample_predictions = np.array(sample_predictions)\n            uncertainty = np.std(sample_predictions, axis=0)\n            \n            return point_prediction, uncertainty\n        \n        return point_prediction\n    \n    def get_feature_importance(self) -> Dict[str, float]:\n        \"\"\"Get posterior weights as importance\"\"\"\n        return dict(zip(self.posterior_weights.keys() if hasattr(self.posterior_weights, 'keys') \n                       else range(len(self.posterior_weights)), self.posterior_weights))\n\n# =============================================================================\n# Neural Attention Meta-Learner\n# =============================================================================\n\nclass NeuralAttentionMetaLearner(BaseMetaLearner):\n    \"\"\"Neural network with attention mechanism for model weighting\"\"\"\n    \n    def __init__(self, hidden_size: int = 64, n_heads: int = 4, dropout: float = 0.3, \n                 learning_rate: float = 0.001, n_epochs: int = 100, random_state: int = 42):\n        self.hidden_size = hidden_size\n        self.n_heads = n_heads\n        self.dropout = dropout\n        self.learning_rate = learning_rate\n        self.n_epochs = n_epochs\n        self.random_state = random_state\n        self.model = None\n        self.scaler = StandardScaler()\n        self.attention_weights_history = []\n    \n    def _build_attention_model(self, n_models: int):\n        \"\"\"Build attention-based neural network using sklearn\"\"\"\n        # Simplified attention mechanism using MLPRegressor\n        # In production, you'd use PyTorch/TensorFlow for true attention\n        \n        self.model = MLPRegressor(\n            hidden_layer_sizes=(self.hidden_size * 2, self.hidden_size, self.hidden_size // 2),\n            activation='relu',\n            solver='adam',\n            alpha=0.01,\n            learning_rate_init=self.learning_rate,\n            max_iter=self.n_epochs,\n            random_state=self.random_state,\n            early_stopping=True,\n            validation_fraction=0.15,\n            n_iter_no_change=10\n        )\n    \n    def _compute_attention_weights(self, predictions: np.ndarray) -> np.ndarray:\n        \"\"\"Compute pseudo-attention weights\"\"\"\n        # Simplified attention: use correlations between predictions\n        n_samples, n_models = predictions.shape\n        \n        # Compute pairwise similarities\n        attention_scores = np.zeros((n_samples, n_models))\n        \n        for i in range(n_models):\n            for j in range(n_models):\n                if i != j:\n                    # Use rolling correlation as attention score\n                    window = min(50, n_samples // 10)\n                    for k in range(window, n_samples):\n                        if k < window:\n                            corr = 0.5\n                        else:\n                            corr = np.corrcoef(\n                                predictions[k-window:k, i],\n                                predictions[k-window:k, j]\n                            )[0, 1]\n                            if np.isnan(corr):\n                                corr = 0.5\n                        attention_scores[k, i] += corr\n        \n        # Normalize to create attention weights\n        attention_weights = np.exp(attention_scores)\n        attention_weights = attention_weights / (np.sum(attention_weights, axis=1, keepdims=True) + 1e-8)\n        \n        return attention_weights\n    \n    def fit(self, predictions: Dict[str, np.ndarray], y: np.ndarray, **kwargs):\n        \"\"\"Fit neural attention model\"\"\"\n        print(\"\\nTraining Neural Attention Meta-Learner...\")\n        \n        # Prepare data\n        pred_matrix = np.column_stack(list(predictions.values()))\n        n_models = pred_matrix.shape[1]\n        \n        # Build model\n        self._build_attention_model(n_models)\n        \n        # Create attention-weighted features\n        attention_weights = self._compute_attention_weights(pred_matrix)\n        self.attention_weights_history = attention_weights\n        \n        # Create enhanced features\n        features = []\n        features.append(pred_matrix)  # Original predictions\n        features.append(pred_matrix ** 2)  # Squared predictions\n        features.append(attention_weights)  # Attention weights\n        \n        # Weighted predictions\n        weighted_preds = pred_matrix * attention_weights\n        features.append(weighted_preds)\n        \n        # Concatenate all features\n        X = np.hstack(features)\n        X_scaled = self.scaler.fit_transform(X)\n        \n        # Train model\n        self.model.fit(X_scaled, y)\n        \n        # Evaluate\n        y_pred = self.model.predict(X_scaled)\n        score = pearsonr(y_pred, y)[0]\n        print(f\"  Training correlation: {score:.4f}\")\n        \n        # Analyze attention patterns\n        mean_attention = np.mean(attention_weights, axis=0)\n        print(\"  Average attention weights:\")\n        for name, weight in zip(predictions.keys(), mean_attention):\n            print(f\"    {name}: {weight:.3f}\")\n    \n    def predict(self, predictions: Dict[str, np.ndarray], **kwargs) -> np.ndarray:\n        \"\"\"Predict using attention model\"\"\"\n        # Prepare data\n        pred_matrix = np.column_stack(list(predictions.values()))\n        \n        # Compute attention weights\n        attention_weights = self._compute_attention_weights(pred_matrix)\n        \n        # Create features\n        features = []\n        features.append(pred_matrix)\n        features.append(pred_matrix ** 2)\n        features.append(attention_weights)\n        features.append(pred_matrix * attention_weights)\n        \n        X = np.hstack(features)\n        X_scaled = self.scaler.transform(X)\n        \n        return self.model.predict(X_scaled)\n    \n    def get_feature_importance(self) -> Dict[str, float]:\n        \"\"\"Get average attention weights as importance\"\"\"\n        if len(self.attention_weights_history) > 0:\n            mean_attention = np.mean(self.attention_weights_history, axis=0)\n            return dict(enumerate(mean_attention))\n        return {}\n\n# =============================================================================\n# Advanced Meta-Learning Ensemble Builder\n# =============================================================================\n\nclass AdvancedMetaLearningEnsemble:\n    \"\"\"Main ensemble builder using multiple meta-learning strategies\"\"\"\n    \n    def __init__(self, strategies: List[str] = None):\n        # Base configuration\n        self.model_name = \"meta_ensemble\"\n        self.model_directory = os.path.join(global_config.base_dir, \"meta_ensemble\")\n        \n        # Register with global configuration\n        global_config.register_model(self.model_name, self.model_directory)\n        \n        # Output paths\n        self.final_submission_path = \"/kaggle/working/final_submission.csv\"\n        self.ensemble_analysis_path = os.path.join(self.model_directory, \"meta_ensemble_analysis.csv\")\n        self.model_correlations_path = os.path.join(self.model_directory, \"model_correlations.csv\")\n        self.meta_weights_path = os.path.join(self.model_directory, \"meta_weights.json\")\n        self.visualization_path = os.path.join(self.model_directory, \"meta_ensemble_visualization.png\")\n        self.meta_models_path = os.path.join(self.model_directory, \"meta_models.pkl\")\n        \n        # Meta-learning strategies\n        if strategies is None:\n            strategies = ['stacking', 'dynamic', 'bayesian', 'attention']\n        self.strategies = strategies\n        \n        # Initialize meta-learners\n        self.meta_learners = {\n            'stacking': MultiLevelStackingMetaLearner(n_levels=2),\n            'dynamic': DynamicModelSelector(n_regimes=4),\n            'bayesian': BayesianModelAveraging(prior_strength=1.0),\n            'attention': NeuralAttentionMetaLearner(hidden_size=64)\n        }\n        \n        # Ensemble parameters\n        self.n_folds = 5\n        self.random_state = 42\n        self.final_blender = None\n        \n        # Ensure directories exist\n        Path(self.model_directory).mkdir(parents=True, exist_ok=True)\n    \n    def load_base_predictions(self) -> Tuple[Dict[str, pd.DataFrame], Dict[str, np.ndarray]]:\n        \"\"\"Load predictions from base models\"\"\"\n        submissions = {}\n        predictions = {}\n        \n        print(\"\\nLoading base model predictions...\")\n        \n        # Define target models\n        target_models = {\n            'xgboost': ['submission', 'xgboost_submission'],\n            'autoencoder_simple': ['ensemble_simple_submission', 'autoencoder_simple'],\n            'autoencoder_weighted': ['ensemble_weighted_submission', 'autoencoder_weighted']\n        }\n        \n        # Load from registry\n        for model_name, possible_output_names in target_models.items():\n            if model_name in global_config.model_registry:\n                model_record = global_config.model_registry[model_name]\n                \n                if model_record.status == 'completed':\n                    for output_name in possible_output_names:\n                        if output_name in model_record.outputs:\n                            output = model_record.outputs[output_name]\n                            if os.path.exists(output.file_path):\n                                try:\n                                    submission_df = pd.read_csv(output.file_path)\n                                    submissions[model_name] = submission_df\n                                    predictions[model_name] = submission_df.iloc[:, 1].values\n                                    print(f\"  ✓ Loaded {model_name}\")\n                                    break\n                                except Exception as e:\n                                    print(f\"  ✗ Error loading {model_name}: {e}\")\n        \n        # Fallback paths\n        fallback_paths = {\n            'xgboost': os.path.join(global_config.base_dir, \"triple_xgboost\", \"submission.csv\"),\n            'autoencoder_simple': os.path.join(global_config.base_dir, \"autoencoder_deepmlp\", \"ensemble_simple_submission.csv\"),\n            'autoencoder_weighted': os.path.join(global_config.base_dir, \"autoencoder_deepmlp\", \"ensemble_weighted_submission.csv\")\n        }\n        \n        for model_name, path in fallback_paths.items():\n            if model_name not in submissions and os.path.exists(path):\n                try:\n                    submission_df = pd.read_csv(path)\n                    submissions[model_name] = submission_df\n                    predictions[model_name] = submission_df.iloc[:, 1].values\n                    print(f\"  ✓ Loaded {model_name} (fallback)\")\n                except Exception as e:\n                    print(f\"  ✗ Error: {e}\")\n        \n        print(f\"\\nLoaded {len(predictions)} models\")\n        return submissions, predictions\n    \n    def create_synthetic_labels(self, predictions: Dict[str, np.ndarray]) -> np.ndarray:\n        \"\"\"Create synthetic labels for meta-learning\"\"\"\n        # Use weighted average based on model diversity\n        pred_array = np.column_stack(list(predictions.values()))\n        \n        # Calculate pairwise correlations\n        n_models = pred_array.shape[1]\n        correlations = np.corrcoef(pred_array.T)\n        \n        # Weight based on uniqueness (lower correlation with others)\n        uniqueness = 1 - (np.sum(correlations, axis=1) - 1) / (n_models - 1)\n        weights = uniqueness / np.sum(uniqueness)\n        \n        # Create synthetic target\n        synthetic_y = np.average(pred_array, axis=1, weights=weights)\n        \n        # Add small noise for regularization\n        np.random.seed(self.random_state)\n        noise = np.random.normal(0, np.std(synthetic_y) * 0.01, size=len(synthetic_y))\n        synthetic_y += noise\n        \n        return synthetic_y\n    \n    def train_meta_learners(self, predictions: Dict[str, np.ndarray], y: np.ndarray) -> Dict[str, float]:\n        \"\"\"Train all meta-learners and evaluate performance\"\"\"\n        scores = {}\n        \n        print(\"\\nTraining meta-learners...\")\n        \n        # Split data for validation\n        split_idx = int(len(y) * 0.8)\n        train_predictions = {k: v[:split_idx] for k, v in predictions.items()}\n        val_predictions = {k: v[split_idx:] for k, v in predictions.items()}\n        y_train = y[:split_idx]\n        y_val = y[split_idx:]\n        \n        # Train each meta-learner\n        for strategy in self.strategies:\n            if strategy in self.meta_learners:\n                print(f\"\\n--- {strategy.upper()} META-LEARNER ---\")\n                \n                try:\n                    # Train\n                    self.meta_learners[strategy].fit(train_predictions, y_train)\n                    \n                    # Validate\n                    val_pred = self.meta_learners[strategy].predict(val_predictions)\n                    score = pearsonr(val_pred, y_val)[0]\n                    scores[strategy] = score\n                    \n                    print(f\"  Validation score: {score:.4f}\")\n                    \n                except Exception as e:\n                    print(f\"  Error: {e}\")\n                    scores[strategy] = 0.0\n        \n        return scores\n    \n    def create_final_ensemble(self, predictions: Dict[str, np.ndarray], \n                            meta_scores: Dict[str, float]) -> np.ndarray:\n        \"\"\"Create final ensemble from meta-learners\"\"\"\n        print(\"\\nCreating final ensemble...\")\n        \n        # Get predictions from each meta-learner\n        meta_predictions = {}\n        \n        for strategy in self.strategies:\n            if strategy in self.meta_learners and meta_scores.get(strategy, 0) > 0:\n                try:\n                    pred = self.meta_learners[strategy].predict(predictions)\n                    meta_predictions[strategy] = pred\n                except:\n                    pass\n        \n        if not meta_predictions:\n            # Fallback to simple average\n            print(\"  Warning: No meta-learners succeeded, using simple average\")\n            return np.mean(list(predictions.values()), axis=0)\n        \n        # Weight meta-learners by performance\n        weights = np.array([meta_scores.get(name, 0) for name in meta_predictions.keys()])\n        weights = np.maximum(weights, 0)  # Ensure non-negative\n        \n        if np.sum(weights) > 0:\n            weights = weights / np.sum(weights)\n        else:\n            weights = np.ones(len(weights)) / len(weights)\n        \n        print(\"\\n  Meta-learner weights:\")\n        for name, weight in zip(meta_predictions.keys(), weights):\n            print(f\"    {name}: {weight:.3f}\")\n        \n        # Create weighted ensemble\n        final_predictions = np.average(\n            list(meta_predictions.values()),\n            axis=0,\n            weights=weights\n        )\n        \n        return final_predictions\n    \n    def save_meta_models(self):\n        \"\"\"Save trained meta-learners\"\"\"\n        with open(self.meta_models_path, 'wb') as f:\n            pickle.dump({\n                'meta_learners': self.meta_learners,\n                'final_blender': self.final_blender\n            }, f)\n        \n        global_config.register_model_output(\n            self.model_name,\n            'meta_models',\n            self.meta_models_path,\n            'model'\n        )\n    \n    def create_advanced_visualization(self, predictions: Dict[str, np.ndarray],\n                                    final_predictions: np.ndarray,\n                                    meta_scores: Dict[str, float]):\n        \"\"\"Create comprehensive visualization of meta-learning results\"\"\"\n        try:\n            fig = plt.figure(figsize=(16, 12))\n            \n            # Create grid\n            gs = fig.add_gridspec(3, 3, hspace=0.3, wspace=0.3)\n            \n            # 1. Base model predictions distribution\n            ax1 = fig.add_subplot(gs[0, :2])\n            for name, pred in predictions.items():\n                ax1.hist(pred, bins=50, alpha=0.5, label=name, density=True)\n            ax1.set_xlabel('Prediction Value')\n            ax1.set_ylabel('Density')\n            ax1.set_title('Base Model Prediction Distributions')\n            ax1.legend()\n            ax1.grid(True, alpha=0.3)\n            \n            # 2. Meta-learner performance\n            ax2 = fig.add_subplot(gs[0, 2])\n            meta_names = list(meta_scores.keys())\n            meta_values = list(meta_scores.values())\n            bars = ax2.bar(meta_names, meta_values, color=['#1f77b4', '#ff7f0e', '#2ca02c', '#d62728'])\n            ax2.set_ylabel('Validation Score')\n            ax2.set_title('Meta-Learner Performance')\n            ax2.set_ylim(0, 1)\n            for i, (name, value) in enumerate(zip(meta_names, meta_values)):\n                ax2.text(i, value + 0.01, f'{value:.3f}', ha='center', va='bottom')\n            \n            # 3. Model correlation heatmap\n            ax3 = fig.add_subplot(gs[1, 0])\n            pred_array = np.column_stack(list(predictions.values()))\n            corr_matrix = np.corrcoef(pred_array.T)\n            im = ax3.imshow(corr_matrix, cmap='coolwarm', vmin=-1, vmax=1)\n            ax3.set_xticks(range(len(predictions)))\n            ax3.set_yticks(range(len(predictions)))\n            ax3.set_xticklabels(list(predictions.keys()), rotation=45, ha='right')\n            ax3.set_yticklabels(list(predictions.keys()))\n            ax3.set_title('Model Correlations')\n            plt.colorbar(im, ax=ax3)\n            \n            # 4. Final vs base predictions scatter\n            ax4 = fig.add_subplot(gs[1, 1])\n            mean_base = np.mean(list(predictions.values()), axis=0)\n            ax4.scatter(mean_base, final_predictions, alpha=0.5, s=1)\n            ax4.plot([mean_base.min(), mean_base.max()], \n                    [mean_base.min(), mean_base.max()], \n                    'r--', lw=2)\n            ax4.set_xlabel('Base Ensemble (Simple Average)')\n            ax4.set_ylabel('Meta-Learning Ensemble')\n            ax4.set_title('Meta-Learning vs Simple Ensemble')\n            ax4.grid(True, alpha=0.3)\n            \n            # 5. Feature importance from stacking\n            ax5 = fig.add_subplot(gs[1, 2])\n            if 'stacking' in self.meta_learners:\n                importance = self.meta_learners['stacking'].get_feature_importance()\n                if importance:\n                    # Show top 10 features\n                    sorted_features = sorted(importance.items(), key=lambda x: x[1], reverse=True)[:10]\n                    feature_names = [f[0] for f in sorted_features]\n                    feature_values = [f[1] for f in sorted_features]\n                    \n                    ax5.barh(range(len(feature_names)), feature_values)\n                    ax5.set_yticks(range(len(feature_names)))\n                    ax5.set_yticklabels(feature_names)\n                    ax5.set_xlabel('Importance')\n                    ax5.set_title('Top 10 Stacking Features')\n            \n            # 6. Final prediction distribution\n            ax6 = fig.add_subplot(gs[2, :])\n            ax6.hist(final_predictions, bins=100, color='darkgreen', alpha=0.8, edgecolor='black')\n            ax6.axvline(np.mean(final_predictions), color='red', linestyle='--', \n                       label=f'Mean: {np.mean(final_predictions):.6f}')\n            ax6.axvline(np.median(final_predictions), color='orange', linestyle='--', \n                       label=f'Median: {np.median(final_predictions):.6f}')\n            ax6.set_xlabel('Prediction Value')\n            ax6.set_ylabel('Count')\n            ax6.set_title('Final Meta-Ensemble Predictions')\n            ax6.legend()\n            ax6.grid(True, alpha=0.3)\n            \n            # Add text summary\n            summary_text = f\"\"\"\nMeta-Learning Ensemble Summary:\n- Models used: {len(predictions)}\n- Meta-learners: {', '.join(self.strategies)}\n- Best meta-learner: {max(meta_scores.items(), key=lambda x: x[1])[0] if meta_scores else 'N/A'}\n- Prediction range: [{np.min(final_predictions):.6f}, {np.max(final_predictions):.6f}]\n- Standard deviation: {np.std(final_predictions):.6f}\n            \"\"\"\n            fig.text(0.02, 0.02, summary_text, fontsize=10, \n                    bbox=dict(boxstyle=\"round,pad=0.5\", facecolor=\"lightgray\"))\n            \n            plt.suptitle('DRW Crypto Prediction - Meta-Learning Ensemble Analysis', fontsize=16)\n            plt.tight_layout()\n            plt.savefig(self.visualization_path, dpi=150, bbox_inches='tight')\n            plt.close()\n            \n            print(f\"\\nVisualization saved to: {self.visualization_path}\")\n            \n            global_config.register_model_output(\n                self.model_name,\n                'visualization',\n                self.visualization_path,\n                'visualization'\n            )\n            \n        except Exception as e:\n            print(f\"Warning: Could not create visualization: {e}\")\n    \n    def run(self) -> pd.DataFrame:\n        \"\"\"Execute the complete meta-learning ensemble pipeline\"\"\"\n        print(\"\\n\" + \"=\"*80)\n        print(\"ADVANCED META-LEARNING ENSEMBLE BUILDER\")\n        print(\"=\"*80)\n        \n        # Update status\n        global_config.update_model_status(self.model_name, 'running')\n        \n        try:\n            # Load base predictions\n            submissions, predictions = self.load_base_predictions()\n            \n            if len(predictions) < 2:\n                raise ValueError(f\"Insufficient models for ensemble: {len(predictions)}\")\n            \n            # Create synthetic labels\n            synthetic_y = self.create_synthetic_labels(predictions)\n            \n            # Train meta-learners\n            meta_scores = self.train_meta_learners(predictions, synthetic_y)\n            \n            # Create final ensemble\n            final_predictions = self.create_final_ensemble(predictions, meta_scores)\n            \n            # Post-processing\n            final_predictions = self.apply_post_processing(final_predictions, predictions)\n            \n            # Create submission\n            template = next(iter(submissions.values()))\n            final_submission = template.copy()\n            final_submission.iloc[:, 1] = final_predictions\n            \n            # Save outputs\n            final_submission.to_csv(self.final_submission_path, index=False)\n            print(f\"\\nFinal submission saved to: {self.final_submission_path}\")\n            \n            global_config.register_model_output(\n                self.model_name,\n                'final_submission',\n                self.final_submission_path,\n                'submission'\n            )\n            \n            # Save meta-learner weights\n            meta_info = {\n                'strategies_used': self.strategies,\n                'meta_scores': meta_scores,\n                'n_base_models': len(predictions),\n                'base_models': list(predictions.keys()),\n                'prediction_stats': {\n                    'mean': float(np.mean(final_predictions)),\n                    'std': float(np.std(final_predictions)),\n                    'min': float(np.min(final_predictions)),\n                    'max': float(np.max(final_predictions))\n                }\n            }\n            \n            with open(self.meta_weights_path, 'w') as f:\n                json.dump(meta_info, f, indent=2)\n            \n            global_config.register_model_output(\n                self.model_name,\n                'meta_weights',\n                self.meta_weights_path,\n                'config',\n                metadata=meta_info\n            )\n            \n            # Save models\n            self.save_meta_models()\n            \n            # Create visualization\n            self.create_advanced_visualization(predictions, final_predictions, meta_scores)\n            \n            # Update status\n            best_score = max(meta_scores.values()) if meta_scores else 0.0\n            global_config.update_model_status(self.model_name, 'completed', score=best_score)\n            \n            print(\"\\n✅ Meta-learning ensemble completed successfully\")\n            \n            return final_submission\n            \n        except Exception as e:\n            error_msg = str(e)\n            global_config.update_model_status(self.model_name, 'failed', error_message=error_msg)\n            print(f\"\\n❌ Meta-learning ensemble failed: {error_msg}\")\n            raise\n    \n    def apply_post_processing(self, predictions: np.ndarray, \n                            base_predictions: Dict[str, np.ndarray]) -> np.ndarray:\n        \"\"\"Apply sophisticated post-processing\"\"\"\n        # Calculate bounds from base predictions\n        all_base = np.concatenate(list(base_predictions.values()))\n        \n        # Use robust percentiles\n        lower_bound = np.percentile(all_base, 0.1)\n        upper_bound = np.percentile(all_base, 99.9)\n        \n        # Clip predictions\n        clipped = np.clip(predictions, lower_bound, upper_bound)\n        \n        # Smooth extreme outliers\n        z_scores = np.abs((clipped - np.median(clipped)) / (1.4826 * np.median(np.abs(clipped - np.median(clipped)))))\n        outlier_mask = z_scores > 3.5\n        \n        if np.sum(outlier_mask) > 0:\n            print(f\"\\nPost-processing {np.sum(outlier_mask)} outliers\")\n            \n            # Use robust location estimate\n            robust_center = np.median(clipped[~outlier_mask])\n            clipped[outlier_mask] = 0.7 * robust_center + 0.3 * clipped[outlier_mask]\n        \n        return clipped\n\n# =============================================================================\n# Main Execution\n# =============================================================================\n\ndef create_meta_learning_ensemble():\n    \"\"\"Create the final ensemble using meta-learning techniques\"\"\"\n    print(\"\\nDRW Crypto Market Prediction - Meta-Learning Ensemble\")\n    print(\"=\"*80)\n    \n    try:\n        # Clean memory\n        aggressive_memory_cleanup()\n        \n        # Choose strategies\n        strategies = ['stacking', 'dynamic', 'bayesian', 'attention']\n        \n        # Create ensemble\n        ensemble = AdvancedMetaLearningEnsemble(strategies=strategies)\n        \n        # Run pipeline\n        final_submission = ensemble.run()\n        \n        return final_submission\n        \n    except Exception as e:\n        print(f\"\\nError: {e}\")\n        raise\n\n# Entry point\nif __name__ == \"__main__\":\n    final_submission = create_meta_learning_ensemble()\n    \n    # Display summary\n    print(\"\\n\" + \"=\"*80)\n    print(\"Pipeline Execution Summary:\")\n    print(global_config.get_execution_summary())\n    print(\"=\"*80)","metadata":{"trusted":true},"outputs":[],"execution_count":null},{"cell_type":"code","source":"# # Final Ensemble Building\n# # !/usr/bin/env python\n# # -*- coding: utf-8 -*-\n# \"\"\"\n# DRW Crypto Market Prediction - Final Ensemble Builder\n# This module creates the final ensemble from previously executed individual model predictions\n# Currently configured to work with XGBoost and AutoEncoder models\n# \"\"\"\n\n# import subprocess\n# import sys\n# import os\n# import gc\n# import warnings\n# import json\n# import pandas as pd\n# import numpy as np\n# from scipy.stats import pearsonr\n# from typing import List, Dict, Tuple, Optional, Any\n# from pathlib import Path\n# from sklearn.model_selection import KFold, GridSearchCV\n# from sklearn.linear_model import Ridge\n# import matplotlib.pyplot as plt\n# import seaborn as sns\n\n# warnings.filterwarnings('ignore')\n\n# # Install required packages for ensemble building\n# print(\"Installing packages for Final Ensemble Builder...\")\n# packages_to_install = [\n#     'flaml==2.1.1',\n#     'matplotlib',\n#     'seaborn'\n# ]\n\n# for package in packages_to_install:\n#     subprocess.check_call([sys.executable, \"-m\", \"pip\", \"install\", package, \"--quiet\"])\n\n# class AdvancedFinalEnsembleBuilder:\n#     \"\"\"Advanced ensemble builder with stacking, cross-validation, and AutoML options\"\"\"\n    \n#     def __init__(self, use_automl: bool = True):\n#         self.use_automl = use_automl\n        \n#         # Model configuration\n#         self.model_name = \"final_ensemble\"\n#         self.model_directory = os.path.join(global_config.base_dir, \"final_ensemble\")\n        \n#         # Register with global configuration\n#         global_config.register_model(self.model_name, self.model_directory)\n        \n#         # Output paths\n#         self.final_submission_path = \"/kaggle/working/final_submission.csv\"\n#         self.ensemble_analysis_path = os.path.join(self.model_directory, \"ensemble_analysis.csv\")\n#         self.model_correlations_path = os.path.join(self.model_directory, \"model_correlations.csv\")\n#         self.ensemble_weights_path = os.path.join(self.model_directory, \"ensemble_weights.json\")\n#         self.visualization_path = os.path.join(self.model_directory, \"ensemble_analysis.png\")\n        \n#         # Ensemble parameters\n#         self.n_folds = 5\n#         self.random_state = 42\n#         self.test_size = 0.2\n        \n#         # Ensure directories exist\n#         Path(self.model_directory).mkdir(parents=True, exist_ok=True)\n    \n#     def load_available_submissions(self) -> Tuple[Dict[str, pd.DataFrame], Dict[str, np.ndarray]]:\n#         \"\"\"Load submissions from successfully completed models using the global registry\"\"\"\n#         submissions = {}\n#         predictions = {}\n        \n#         print(\"\\nSearching for model submissions in global registry...\")\n        \n#         # Define target models and their expected submission names\n#         target_models = {\n#             'xgboost': ['submission', 'xgboost_submission'],\n#             'autoencoder_simple': ['ensemble_simple_submission', 'autoencoder_simple'],\n#             'autoencoder_weighted': ['ensemble_weighted_submission', 'autoencoder_weighted']\n#         }\n        \n#         # First attempt: Check the global registry for registered outputs\n#         for model_name, possible_output_names in target_models.items():\n#             if model_name in global_config.model_registry:\n#                 model_record = global_config.model_registry[model_name]\n                \n#                 # Check if model completed successfully\n#                 if model_record.status == 'completed':\n#                     # Look for submission files in the model's outputs\n#                     submission_found = False\n#                     for output_name in possible_output_names:\n#                         if output_name in model_record.outputs:\n#                             output = model_record.outputs[output_name]\n#                             if os.path.exists(output.file_path):\n#                                 try:\n#                                     submission_df = pd.read_csv(output.file_path)\n#                                     submissions[model_name] = submission_df\n#                                     predictions[model_name] = submission_df.iloc[:, 1].values\n#                                     print(f\"  ✓ Loaded {model_name} from registry: {output.file_path}\")\n#                                     submission_found = True\n#                                     break\n#                                 except Exception as e:\n#                                     print(f\"  ✗ Error loading {model_name} from {output.file_path}: {e}\")\n                    \n#                     if not submission_found:\n#                         print(f\"  ⚠ {model_name} completed but no submission found in outputs\")\n#                 else:\n#                     print(f\"  ⚠ {model_name} status: {model_record.status}\")\n#             else:\n#                 print(f\"  ⚠ {model_name} not found in registry\")\n        \n#         # Fallback: Check standard file paths if registry is incomplete\n#         fallback_paths = {\n#             'xgboost': os.path.join(global_config.base_dir, \"triple_xgboost\", \"submission.csv\"),\n#             'autoencoder_simple': os.path.join(global_config.base_dir, \"autoencoder_deepmlp\", \"ensemble_simple_submission.csv\"),\n#             'autoencoder_weighted': os.path.join(global_config.base_dir, \"autoencoder_deepmlp\", \"ensemble_weighted_submission.csv\")\n#         }\n        \n#         for model_name, submission_path in fallback_paths.items():\n#             if model_name not in submissions and os.path.exists(submission_path):\n#                 try:\n#                     submission_df = pd.read_csv(submission_path)\n#                     submissions[model_name] = submission_df\n#                     predictions[model_name] = submission_df.iloc[:, 1].values\n#                     print(f\"  ✓ Loaded {model_name} from fallback path: {submission_path}\")\n#                 except Exception as e:\n#                     print(f\"  ✗ Error loading {model_name} from fallback: {e}\")\n        \n#         # Validate loaded predictions\n#         print(f\"\\nLoaded {len(submissions)} model submissions:\")\n#         for model_name, preds in predictions.items():\n#             print(f\"  • {model_name}: {len(preds)} predictions, range [{preds.min():.4f}, {preds.max():.4f}]\")\n        \n#         return submissions, predictions\n    \n#     def analyze_model_diversity(self, predictions: Dict[str, np.ndarray]) -> pd.DataFrame:\n#         \"\"\"Analyze correlation between model predictions to ensure diversity\"\"\"\n#         model_names = list(predictions.keys())\n#         n_models = len(model_names)\n        \n#         if n_models < 2:\n#             print(\"Warning: Insufficient models for diversity analysis\")\n#             return pd.DataFrame()\n        \n#         correlation_matrix = np.zeros((n_models, n_models))\n        \n#         for i, model1 in enumerate(model_names):\n#             for j, model2 in enumerate(model_names):\n#                 if i <= j:\n#                     corr = pearsonr(predictions[model1], predictions[model2])[0]\n#                     correlation_matrix[i, j] = corr\n#                     correlation_matrix[j, i] = corr\n        \n#         corr_df = pd.DataFrame(\n#             correlation_matrix,\n#             index=model_names,\n#             columns=model_names\n#         )\n        \n#         # Save correlation matrix\n#         corr_df.to_csv(self.model_correlations_path)\n#         global_config.register_model_output(\n#             self.model_name, \n#             'model_correlations', \n#             self.model_correlations_path,\n#             'analysis'\n#         )\n        \n#         # Display correlation matrix\n#         print(\"\\nModel Correlation Matrix:\")\n#         print(corr_df.round(4))\n        \n#         # Calculate diversity metrics\n#         off_diagonal_corr = correlation_matrix[np.triu_indices(n_models, k=1)]\n        \n#         print(f\"\\nDiversity Metrics:\")\n#         print(f\"  Average inter-model correlation: {np.mean(off_diagonal_corr):.4f}\")\n#         if len(off_diagonal_corr) > 0:\n#             print(f\"  Correlation range: [{np.min(off_diagonal_corr):.4f}, {np.max(off_diagonal_corr):.4f}]\")\n        \n#         return corr_df\n    \n#     def create_stacking_features(self, predictions: Dict[str, np.ndarray]) -> np.ndarray:\n#         \"\"\"Create comprehensive feature matrix for stacking ensemble\"\"\"\n#         feature_list = []\n#         feature_names = []\n        \n#         # Add base predictions\n#         for model_name, preds in predictions.items():\n#             feature_list.append(preds)\n#             feature_names.append(model_name)\n        \n#         # Add interaction features between different model types\n#         model_names = list(predictions.keys())\n#         for i in range(len(model_names)):\n#             for j in range(i + 1, len(model_names)):\n#                 interaction = predictions[model_names[i]] * predictions[model_names[j]]\n#                 feature_list.append(interaction)\n#                 feature_names.append(f\"{model_names[i]}*{model_names[j]}\")\n        \n#         # Statistical features across all predictions\n#         predictions_array = np.column_stack(list(predictions.values()))\n        \n#         # Basic statistics\n#         feature_list.extend([\n#             np.mean(predictions_array, axis=1),\n#             np.std(predictions_array, axis=1),\n#             np.min(predictions_array, axis=1),\n#             np.max(predictions_array, axis=1),\n#             np.median(predictions_array, axis=1)\n#         ])\n#         feature_names.extend(['mean', 'std', 'min', 'max', 'median'])\n        \n#         # Advanced statistics\n#         feature_list.extend([\n#             np.max(predictions_array, axis=1) - np.min(predictions_array, axis=1),  # Range\n#             np.where(np.mean(predictions_array, axis=1) != 0, \n#                     np.std(predictions_array, axis=1) / np.abs(np.mean(predictions_array, axis=1)), \n#                     0)  # Coefficient of variation\n#         ])\n#         feature_names.extend(['range', 'cv'])\n        \n#         # Stack all features\n#         feature_matrix = np.column_stack(feature_list)\n        \n#         print(f\"\\nCreated {feature_matrix.shape[1]} stacking features: {', '.join(feature_names)}\")\n        \n#         return feature_matrix\n    \n#     def train_automl_ensemble(self, X: np.ndarray, y: np.ndarray) -> Tuple[Any, float]:\n#         \"\"\"Train ensemble using FLAML AutoML\"\"\"\n#         try:\n#             from flaml import AutoML\n            \n#             print(\"\\nTraining FLAML AutoML ensemble...\")\n            \n#             # Split for validation\n#             split_idx = int(len(X) * 0.8)\n#             X_train, X_val = X[:split_idx], X[split_idx:]\n#             y_train, y_val = y[:split_idx], y[split_idx:]\n            \n#             # Configure AutoML\n#             automl = AutoML()\n            \n#             automl_settings = {\n#                 \"time_budget\": 120,  # 2 minutes\n#                 \"metric\": 'r2',\n#                 \"task\": 'regression',\n#                 \"n_jobs\": -1,\n#                 \"estimator_list\": ['rf', 'xgboost', 'lgbm', 'lrl1', 'lrl2'],\n#                 \"seed\": self.random_state,\n#                 \"verbose\": 0,\n#                 \"eval_method\": \"cv\",\n#                 \"n_splits\": 3\n#             }\n            \n#             # Train AutoML\n#             automl.fit(X_train, y_train, **automl_settings)\n            \n#             # Evaluate\n#             val_pred = automl.predict(X_val)\n#             val_score = pearsonr(y_val, val_pred)[0]\n            \n#             print(f\"  FLAML selected model: {automl.best_estimator}\")\n#             print(f\"  Validation correlation: {val_score:.6f}\")\n            \n#             return automl, val_score\n            \n#         except Exception as e:\n#             print(f\"  AutoML training failed: {e}\")\n#             return None, 0.0\n    \n#     def train_ridge_ensemble(self, X: np.ndarray, y: np.ndarray) -> Tuple[Ridge, float]:\n#         \"\"\"Train Ridge regression ensemble with hyperparameter optimization\"\"\"\n#         print(\"\\nTraining Ridge regression ensemble...\")\n        \n#         # Split for validation\n#         split_idx = int(len(X) * 0.8)\n#         X_train, X_val = X[:split_idx], X[split_idx:]\n#         y_train, y_val = y[:split_idx], y[split_idx:]\n        \n#         # Test multiple alpha values\n#         alphas = [0.001, 0.01, 0.1, 1.0, 10.0, 100.0]\n#         best_score = -float('inf')\n#         best_model = None\n#         best_alpha = None\n        \n#         for alpha in alphas:\n#             model = Ridge(alpha=alpha, random_state=self.random_state)\n#             model.fit(X_train, y_train)\n            \n#             val_pred = model.predict(X_val)\n#             score = pearsonr(y_val, val_pred)[0]\n            \n#             if score > best_score:\n#                 best_score = score\n#                 best_model = model\n#                 best_alpha = alpha\n        \n#         print(f\"  Best alpha: {best_alpha}\")\n#         print(f\"  Validation correlation: {best_score:.6f}\")\n        \n#         return best_model, best_score\n    \n#     def create_ensemble(self, predictions: Dict[str, np.ndarray]) -> np.ndarray:\n#         \"\"\"Create the final ensemble using the best available method\"\"\"\n#         # Create feature matrix\n#         X = self.create_stacking_features(predictions)\n        \n#         # Create synthetic target based on weighted average\n#         # Weight based on individual model performance from registry\n#         weights = {}\n#         for model_name in predictions.keys():\n#             if model_name in global_config.model_registry:\n#                 model_score = global_config.model_registry[model_name].score\n#                 if model_score is not None and model_score > 0:\n#                     weights[model_name] = model_score\n#                 else:\n#                     weights[model_name] = 0.1  # Default small weight\n#             else:\n#                 weights[model_name] = 0.1\n        \n#         # Normalize weights\n#         total_weight = sum(weights.values())\n#         weights = {k: v/total_weight for k, v in weights.items()}\n        \n#         print(f\"\\nModel weights based on individual performance:\")\n#         for model_name, weight in weights.items():\n#             print(f\"  • {model_name}: {weight:.3f}\")\n        \n#         # Create weighted synthetic target\n#         y_synthetic = np.zeros(len(next(iter(predictions.values()))))\n#         for model_name, preds in predictions.items():\n#             y_synthetic += weights[model_name] * preds\n        \n#         # Add small noise for training stability\n#         np.random.seed(self.random_state)\n#         noise = np.random.normal(0, np.std(y_synthetic) * 0.02, size=len(y_synthetic))\n#         y_synthetic = y_synthetic + noise\n        \n#         best_model = None\n#         best_score = -float('inf')\n#         best_method = None\n        \n#         # Try AutoML if enabled\n#         if self.use_automl:\n#             automl_model, automl_score = self.train_automl_ensemble(X, y_synthetic)\n#             if automl_model is not None and automl_score > best_score:\n#                 best_model = automl_model\n#                 best_score = automl_score\n#                 best_method = 'AutoML'\n        \n#         # Always try Ridge as baseline\n#         ridge_model, ridge_score = self.train_ridge_ensemble(X, y_synthetic)\n#         if ridge_score > best_score or best_model is None:\n#             best_model = ridge_model\n#             best_score = ridge_score\n#             best_method = 'Ridge'\n        \n#         print(f\"\\nSelected ensemble method: {best_method} (score: {best_score:.6f})\")\n        \n#         # Make final predictions\n#         final_predictions = best_model.predict(X)\n        \n#         # Save ensemble information\n#         ensemble_info = {\n#             'method': best_method,\n#             'validation_score': float(best_score),\n#             'n_models': len(predictions),\n#             'models_used': list(predictions.keys()),\n#             'model_weights': weights\n#         }\n        \n#         with open(self.ensemble_weights_path, 'w') as f:\n#             json.dump(ensemble_info, f, indent=2)\n        \n#         global_config.register_model_output(\n#             self.model_name,\n#             'ensemble_weights',\n#             self.ensemble_weights_path,\n#             'config',\n#             metadata=ensemble_info\n#         )\n        \n#         return final_predictions\n    \n#     def apply_post_processing(self, predictions: np.ndarray, \n#                             base_predictions: Dict[str, np.ndarray]) -> np.ndarray:\n#         \"\"\"Apply post-processing to ensure prediction quality\"\"\"\n#         # Calculate reasonable bounds from base predictions\n#         all_base_preds = np.concatenate(list(base_predictions.values()))\n        \n#         lower_bound = np.percentile(all_base_preds, 0.5)\n#         upper_bound = np.percentile(all_base_preds, 99.5)\n        \n#         # Clip extreme predictions\n#         clipped_predictions = np.clip(predictions, lower_bound, upper_bound)\n        \n#         # Identify and smooth outliers\n#         pred_mean = np.mean(clipped_predictions)\n#         pred_std = np.std(clipped_predictions)\n#         z_scores = np.abs((clipped_predictions - pred_mean) / (pred_std + 1e-8))\n        \n#         outlier_mask = z_scores > 3\n#         n_outliers = np.sum(outlier_mask)\n        \n#         if n_outliers > 0:\n#             print(f\"\\nPost-processing: Adjusting {n_outliers} outliers ({n_outliers/len(predictions)*100:.2f}%)\")\n            \n#             # Blend outliers toward the mean\n#             blend_factor = 0.7  # 70% mean, 30% original\n#             clipped_predictions[outlier_mask] = (\n#                 blend_factor * pred_mean + \n#                 (1 - blend_factor) * clipped_predictions[outlier_mask]\n#             )\n        \n#         return clipped_predictions\n    \n#     def create_visualization(self, final_submission: pd.DataFrame, \n#                            predictions: Dict[str, np.ndarray]):\n#         \"\"\"Create comprehensive visualization of ensemble results\"\"\"\n#         try:\n#             fig, axes = plt.subplots(2, 2, figsize=(14, 10))\n            \n#             # Plot 1: Individual model predictions\n#             ax1 = axes[0, 0]\n#             for model_name, preds in predictions.items():\n#                 ax1.hist(preds, bins=50, alpha=0.6, label=model_name, density=True)\n#             ax1.set_xlabel('Prediction Value')\n#             ax1.set_ylabel('Density')\n#             ax1.set_title('Distribution of Individual Model Predictions')\n#             ax1.legend()\n#             ax1.grid(True, alpha=0.3)\n            \n#             # Plot 2: Final ensemble distribution\n#             ax2 = axes[0, 1]\n#             final_preds = final_submission.iloc[:, 1].values\n#             ax2.hist(final_preds, bins=50, color='darkgreen', alpha=0.8, density=True, edgecolor='black')\n#             ax2.axvline(np.mean(final_preds), color='red', linestyle='--', label=f'Mean: {np.mean(final_preds):.4f}')\n#             ax2.set_xlabel('Prediction Value')\n#             ax2.set_ylabel('Density')\n#             ax2.set_title('Final Ensemble Predictions')\n#             ax2.legend()\n#             ax2.grid(True, alpha=0.3)\n            \n#             # Plot 3: Model agreement\n#             ax3 = axes[1, 0]\n#             predictions_array = np.column_stack(list(predictions.values()))\n#             model_std = np.std(predictions_array, axis=1)\n#             ax3.hist(model_std, bins=50, color='orange', alpha=0.8, edgecolor='black')\n#             ax3.axvline(np.mean(model_std), color='red', linestyle='--', label=f'Mean Std: {np.mean(model_std):.4f}')\n#             ax3.set_xlabel('Standard Deviation Across Models')\n#             ax3.set_ylabel('Count')\n#             ax3.set_title('Model Agreement Analysis')\n#             ax3.legend()\n#             ax3.grid(True, alpha=0.3)\n            \n#             # Plot 4: Summary statistics\n#             ax4 = axes[1, 1]\n#             ax4.axis('off')\n            \n#             stats_text = f\"\"\"Ensemble Summary Statistics\n            \n# Final Predictions:\n#   • Mean: {np.mean(final_preds):.6f}\n#   • Std Dev: {np.std(final_preds):.6f}\n#   • Range: [{np.min(final_preds):.6f}, {np.max(final_preds):.6f}]\n#   • Median: {np.median(final_preds):.6f}\n\n# Model Agreement:\n#   • Average Std: {np.mean(model_std):.6f}\n#   • Max Disagreement: {np.max(model_std):.6f}\n#   • Models Used: {len(predictions)}\"\"\"\n            \n#             ax4.text(0.1, 0.5, stats_text, fontsize=11, verticalalignment='center', \n#                     fontfamily='monospace', bbox=dict(boxstyle=\"round,pad=0.5\", facecolor=\"lightgray\"))\n            \n#             plt.suptitle('DRW Crypto Prediction - Ensemble Analysis', fontsize=14, fontweight='bold')\n#             plt.tight_layout()\n#             plt.savefig(self.visualization_path, dpi=150, bbox_inches='tight')\n#             plt.close()\n            \n#             print(f\"Visualization saved to: {self.visualization_path}\")\n            \n#             global_config.register_model_output(\n#                 self.model_name,\n#                 'visualization',\n#                 self.visualization_path,\n#                 'visualization'\n#             )\n            \n#         except Exception as e:\n#             print(f\"Warning: Could not create visualization: {e}\")\n    \n#     def run(self) -> pd.DataFrame:\n#         \"\"\"Execute the complete ensemble building process\"\"\"\n#         print(\"\\nBuilding Advanced Ensemble\")\n#         print(\"=\"*80)\n        \n#         # Update status\n#         global_config.update_model_status(self.model_name, 'running')\n        \n#         # Load available submissions\n#         submissions, predictions = self.load_available_submissions()\n        \n#         if len(submissions) < 2:\n#             # Try single model if available\n#             if len(submissions) == 1:\n#                 print(\"\\nWarning: Only one model available, using single model predictions\")\n#                 template = next(iter(submissions.values()))\n#                 final_submission = template.copy()\n#                 final_submission.to_csv(self.final_submission_path, index=False)\n                \n#                 global_config.register_model_output(\n#                     self.model_name,\n#                     'final_submission',\n#                     self.final_submission_path,\n#                     'submission'\n#                 )\n                \n#                 global_config.update_model_status(self.model_name, 'completed')\n#                 return final_submission\n#             else:\n#                 error_msg = f\"No models available for ensemble\"\n#                 global_config.update_model_status(self.model_name, 'failed', error_message=error_msg)\n#                 raise ValueError(error_msg)\n        \n#         print(f\"\\nSuccessfully loaded {len(submissions)} model submissions for ensemble\")\n        \n#         # Analyze model diversity\n#         correlation_df = self.analyze_model_diversity(predictions)\n        \n#         # Create ensemble\n#         print(\"\\nCreating ensemble predictions...\")\n#         template_submission = next(iter(submissions.values()))\n#         ensemble_predictions = self.create_ensemble(predictions)\n        \n#         # Apply post-processing\n#         final_predictions = self.apply_post_processing(ensemble_predictions, predictions)\n        \n#         # Create final submission\n#         final_submission = template_submission.copy()\n#         final_submission.iloc[:, 1] = final_predictions\n        \n#         # Save results\n#         final_submission.to_csv(self.final_submission_path, index=False)\n#         print(f\"\\nFinal submission saved to: {self.final_submission_path}\")\n        \n#         global_config.register_model_output(\n#             self.model_name,\n#             'final_submission',\n#             self.final_submission_path,\n#             'submission'\n#         )\n        \n#         # Create analysis dataframe\n#         analysis_df = pd.DataFrame({\n#             'final_ensemble': final_predictions,\n#             'pre_postprocess': ensemble_predictions,\n#             **predictions\n#         })\n#         analysis_df.to_csv(self.ensemble_analysis_path, index=False)\n        \n#         global_config.register_model_output(\n#             self.model_name,\n#             'ensemble_analysis',\n#             self.ensemble_analysis_path,\n#             'analysis'\n#         )\n        \n#         # Create visualization\n#         self.create_visualization(final_submission, predictions)\n        \n#         # Update status\n#         global_config.update_model_status(self.model_name, 'completed')\n        \n#         print(\"\\nEnsemble building completed successfully\")\n        \n#         return final_submission\n\n# # Main execution function\n# def create_final_ensemble():\n#     \"\"\"Create the final ensemble from completed model predictions\"\"\"\n#     print(\"\\nDRW Crypto Market Prediction - Final Ensemble Builder\")\n#     print(\"=\"*80)\n    \n#     try:\n#         # Clean memory before starting\n#         aggressive_memory_cleanup()\n        \n#         # Configure ensemble builder\n#         use_automl = True  # Set to False to use only Ridge regression\n#         ensemble_builder = AdvancedFinalEnsembleBuilder(use_automl=use_automl)\n        \n#         # Build ensemble\n#         final_submission = ensemble_builder.run()\n        \n#         print(\"\\n✅ Final ensemble creation completed successfully\")\n        \n#         return final_submission\n        \n#     except Exception as e:\n#         error_msg = str(e)\n#         print(f\"\\n❌ Ensemble creation failed: {error_msg}\")\n#         raise\n\n# # Main entry point\n# if __name__ == \"__main__\":\n#     final_submission = create_final_ensemble()\n    \n#     # Display final execution summary\n#     print(\"\\n\" + \"=\"*80)\n#     print(\"Pipeline Execution Summary:\")\n#     print(global_config.get_execution_summary())\n#     print(\"=\"*80)","metadata":{"trusted":true},"outputs":[],"execution_count":null}]}