{"metadata":{"kernelspec":{"language":"python","display_name":"Python 3","name":"python3"},"language_info":{"pygments_lexer":"ipython3","nbconvert_exporter":"python","version":"3.6.4","file_extension":".py","codemirror_mode":{"name":"ipython","version":3},"name":"python","mimetype":"text/x-python"},"kaggle":{"accelerator":"none","dataSources":[{"sourceId":96164,"databundleVersionId":11418275,"sourceType":"competition"}],"isInternetEnabled":true,"language":"python","sourceType":"notebook","isGpuEnabled":false}},"nbformat_minor":4,"nbformat":4,"cells":[{"cell_type":"code","source":"# This Python 3 environment comes with many helpful analytics libraries installed\n# It is defined by the kaggle/python Docker image: https://github.com/kaggle/docker-python\n# For example, here's several helpful packages to load\n\nimport numpy as np # linear algebra\nimport pandas as pd # data processing, CSV file I/O (e.g. pd.read_csv)\n\n# Input data files are available in the read-only \"../input/\" directory\n# For example, running this (by clicking run or pressing Shift+Enter) will list all files under the input directory\n\nimport os\nfor dirname, _, filenames in os.walk('/kaggle/input'):\n    for filename in filenames:\n        print(os.path.join(dirname, filename))\n\n# You can write up to 20GB to the current directory (/kaggle/working/) that gets preserved as output when you create a version using \"Save & Run All\" \n# You can also write temporary files to /kaggle/temp/, but they won't be saved outside of the current session","metadata":{"_uuid":"8f2839f25d086af736a60e9eeb907d3b93b6e0e5","_cell_guid":"b1076dfc-b9ad-4769-8c92-a6c4dae69d19","trusted":true},"outputs":[],"execution_count":null},{"cell_type":"code","source":"!pip install FLAML\n!pip install ray\n\nimport os\nimport numpy as np\nimport pandas as pd\nfrom sklearn.preprocessing import RobustScaler, PolynomialFeatures, PowerTransformer\nfrom sklearn.feature_selection import VarianceThreshold\nfrom sklearn.cluster import KMeans\nfrom flaml import AutoML\nimport ray\nimport time\nimport gc\nfrom typing import Tuple, Dict, Any\nimport warnings\nimport datetime\n\nwarnings.filterwarnings('ignore')\n\n# Simple timestamp function for logging\ndef get_timestamp():\n    \"\"\"Return current timestamp in the format: YYYY-MM-DD HH:MM:SS\"\"\"\n    return datetime.datetime.now().strftime('%Y-%m-%d %H:%M:%S')\n\n# Print version information\nprint(f\"{get_timestamp()} - INFO - Starting Crypto Price Movement Prediction Pipeline\")\n\nclass CryptoPricePredictor:\n    \"\"\"\n    Crypto market price movement prediction using FLAML AutoML.\n    \n    This implementation focuses on computational efficiency, robust feature\n    engineering, and optimization to deliver high-quality price movement predictions.\n    Enhanced with Ray for distributed hyperparameter tuning.\n    \"\"\"\n    \n    def __init__(self, input_path: str, output_path: str = './', seed: int = 42, num_cpus: int = 4):\n        \"\"\"\n        Initialize the crypto predictor with paths and configuration.\n        \n        Args:\n            input_path: Directory containing train.parquet and test.parquet\n            output_path: Directory to save output files\n            seed: Random seed for reproducibility\n            num_cpus: Number of CPUs to use for Ray distributed computing\n        \"\"\"\n        self.input_path = input_path\n        self.output_path = output_path\n        self.target_column = 'label'\n        self.id_column = 'timestamp'\n        self.start_time = time.time()\n        self.seed = seed\n        self.num_cpus = num_cpus\n        \n        # Key features for crypto market\n        self.market_features = [\n            'bid_qty', 'ask_qty', 'buy_qty', 'sell_qty', 'volume'\n        ]\n        \n        # Anonymous features\n        self.anon_features = [f'X_{i}' for i in range(1, 891)]\n        \n        # Create output directory\n        os.makedirs(self.output_path, exist_ok=True)\n        \n        # Configure time budget (in seconds)\n        self.flaml_time_budget = 1800*8*2.5  # 4 hours\n    \n    def load_data(self) -> Tuple[pd.DataFrame, pd.DataFrame]:\n        \"\"\"\n        Load training and test data with validation and error handling.\n        \n        Returns:\n            Tuple containing training and test dataframes\n        \"\"\"\n        print(f\"{get_timestamp()} - INFO - Loading datasets...\")\n        \n        # Check if files exist\n        train_path = os.path.join(self.input_path, 'train.parquet')\n        test_path = os.path.join(self.input_path, 'test.parquet')\n        \n        if not os.path.exists(train_path):\n            print(f\"{get_timestamp()} - ERROR - Training file not found at {train_path}\")\n            raise FileNotFoundError(f\"Training file not found at {train_path}\")\n        \n        if not os.path.exists(test_path):\n            print(f\"{get_timestamp()} - ERROR - Test file not found at {test_path}\")\n            raise FileNotFoundError(f\"Test file not found at {test_path}\")\n        \n        # Load data\n        train = pd.read_parquet(train_path)\n        test = pd.read_parquet(test_path)\n        \n        # Validate data structure\n        if self.target_column not in train.columns:\n            print(f\"{get_timestamp()} - ERROR - Target column '{self.target_column}' not found in training data\")\n            raise ValueError(f\"Target column '{self.target_column}' not found in training data\")\n        \n        if self.id_column not in test.columns:\n            print(f\"{get_timestamp()} - ERROR - ID column '{self.id_column}' not found in test data\")\n            raise ValueError(f\"ID column '{self.id_column}' not found in test data\")\n        \n        # Log data summary\n        print(f\"{get_timestamp()} - INFO - Loaded datasets - Train: {train.shape}, Test: {test.shape}\")\n        print(f\"{get_timestamp()} - INFO - Training data missing values: {train.isna().sum().sum()} cells\")\n        print(f\"{get_timestamp()} - INFO - Test data missing values: {test.isna().sum().sum()} cells\")\n        \n        # Optimize memory usage with better data types\n        self._optimize_dtypes(train)\n        self._optimize_dtypes(test)\n        \n        return train, test\n    \n    def _optimize_dtypes(self, df: pd.DataFrame) -> None:\n        \"\"\"\n        Optimize data types to reduce memory usage.\n        \n        Args:\n            df: DataFrame to optimize\n        \"\"\"\n        # Handle integer columns\n        int_cols = df.select_dtypes(include=['int']).columns\n        for col in int_cols:\n            # Skip optimization if column has missing values\n            if df[col].isna().any():\n                continue\n                \n            c_min = df[col].min()\n            c_max = df[col].max()\n            \n            # Convert to smallest possible integer type\n            if c_min >= 0:\n                if c_max < 256:\n                    df[col] = df[col].astype(np.uint8)\n                elif c_max < 65536:\n                    df[col] = df[col].astype(np.uint16)\n                else:\n                    df[col] = df[col].astype(np.uint32)\n            else:\n                if c_min > -128 and c_max < 128:\n                    df[col] = df[col].astype(np.int8)\n                elif c_min > -32768 and c_max < 32768:\n                    df[col] = df[col].astype(np.int16)\n                else:\n                    df[col] = df[col].astype(np.int32)\n                    \n        # Handle float columns\n        float_cols = df.select_dtypes(include=['float']).columns\n        for col in float_cols:\n            # Skip optimization if column has missing values\n            if df[col].isna().any():\n                continue\n                \n            df[col] = df[col].astype(np.float32)\n    \n    def engineer_features(self, train: pd.DataFrame, test: pd.DataFrame) -> Tuple[pd.DataFrame, pd.DataFrame]:\n        \"\"\"\n        Apply domain-specific feature engineering for crypto price movement prediction.\n        \n        Args:\n            train: Training DataFrame\n            test: Test DataFrame\n            \n        Returns:\n            Tuple of processed training and test DataFrames\n        \"\"\"\n        print(f\"{get_timestamp()} - INFO - Applying crypto-specific feature engineering...\")\n        \n        # Extract target and IDs before transformations\n        self.y_train_original = train[self.target_column].copy()\n        \n        # Analyze original target distribution\n        self._analyze_target_distribution(self.y_train_original)\n        \n        # Apply winsorization to handle outliers in the target variable\n        upper_percentile = 99.7\n        lower_percentile = 0.3\n        \n        lower_bound = np.percentile(self.y_train_original, lower_percentile)\n        upper_bound = np.percentile(self.y_train_original, upper_percentile)\n        \n        # Apply caps while preserving the order\n        self.y_train_winsorized = np.clip(self.y_train_original, lower_bound, upper_bound)\n        \n        # For crypto, we might want to keep the original scale or use a simpler transformation\n        # Since we're dealing with price movements, let's use a simple standardization\n        self.y_train = self.y_train_winsorized\n        \n        print(f\"{get_timestamp()} - INFO - Applied winsorization to target variable\")\n        print(f\"{get_timestamp()} - INFO - Target stats - original: min={self.y_train_original.min():.4f}, max={self.y_train_original.max():.4f}, mean={self.y_train_original.mean():.4f}\")\n        print(f\"{get_timestamp()} - INFO - Target stats - winsorized: min={self.y_train_winsorized.min():.4f}, max={self.y_train_winsorized.max():.4f}, mean={self.y_train_winsorized.mean():.4f}\")\n        \n        # Store test IDs for later use\n        self.test_ids = test[self.id_column].copy()\n        \n        # Analyze and store column types for dynamic feature creation\n        self.numeric_columns = train.select_dtypes(include=['int', 'float']).columns.tolist()\n        if self.target_column in self.numeric_columns:\n            self.numeric_columns.remove(self.target_column)\n        if self.id_column in self.numeric_columns:\n            self.numeric_columns.remove(self.id_column)\n            \n        print(f\"{get_timestamp()} - INFO - Identified {len(self.numeric_columns)} numeric columns\")\n        \n        # Initialize processed datasets\n        train_processed = train.copy()\n        test_processed = test.copy()\n        \n        # 1. Market microstructure features\n        # Bid-Ask Spread\n        if 'bid_qty' in train.columns and 'ask_qty' in train.columns:\n            # Fill missing values\n            train_processed['bid_qty'] = train_processed['bid_qty'].fillna(0)\n            test_processed['bid_qty'] = test_processed['bid_qty'].fillna(0)\n            train_processed['ask_qty'] = train_processed['ask_qty'].fillna(0)\n            test_processed['ask_qty'] = test_processed['ask_qty'].fillna(0)\n            \n            # Bid-Ask imbalance (order book imbalance)\n            train_processed['bid_ask_imbalance'] = ((train_processed['bid_qty'] - train_processed['ask_qty']) / \n                                                   (train_processed['bid_qty'] + train_processed['ask_qty'] + 1)).astype(np.float32)\n            test_processed['bid_ask_imbalance'] = ((test_processed['bid_qty'] - test_processed['ask_qty']) / \n                                                  (test_processed['bid_qty'] + test_processed['ask_qty'] + 1)).astype(np.float32)\n            \n            # Log transformations\n            train_processed['log_bid_qty'] = np.log1p(train_processed['bid_qty']).astype(np.float32)\n            test_processed['log_bid_qty'] = np.log1p(test_processed['bid_qty']).astype(np.float32)\n            \n            train_processed['log_ask_qty'] = np.log1p(train_processed['ask_qty']).astype(np.float32)\n            test_processed['log_ask_qty'] = np.log1p(test_processed['ask_qty']).astype(np.float32)\n            \n            # Bid-Ask ratio\n            train_processed['bid_ask_ratio'] = (train_processed['bid_qty'] / \n                                               (train_processed['ask_qty'] + 1)).astype(np.float32)\n            test_processed['bid_ask_ratio'] = (test_processed['bid_qty'] / \n                                              (test_processed['ask_qty'] + 1)).astype(np.float32)\n            \n            # Handle infinities\n            train_processed['bid_ask_ratio'] = train_processed['bid_ask_ratio'].replace([np.inf, -np.inf], 0)\n            test_processed['bid_ask_ratio'] = test_processed['bid_ask_ratio'].replace([np.inf, -np.inf], 0)\n        \n        # 2. Trade flow features\n        if 'buy_qty' in train.columns and 'sell_qty' in train.columns:\n            # Fill missing values\n            train_processed['buy_qty'] = train_processed['buy_qty'].fillna(0)\n            test_processed['buy_qty'] = test_processed['buy_qty'].fillna(0)\n            train_processed['sell_qty'] = train_processed['sell_qty'].fillna(0)\n            test_processed['sell_qty'] = test_processed['sell_qty'].fillna(0)\n            \n            # Trade flow imbalance\n            train_processed['trade_flow_imbalance'] = ((train_processed['buy_qty'] - train_processed['sell_qty']) / \n                                                      (train_processed['buy_qty'] + train_processed['sell_qty'] + 1)).astype(np.float32)\n            test_processed['trade_flow_imbalance'] = ((test_processed['buy_qty'] - test_processed['sell_qty']) / \n                                                     (test_processed['buy_qty'] + test_processed['sell_qty'] + 1)).astype(np.float32)\n            \n            # Buy-Sell ratio\n            train_processed['buy_sell_ratio'] = (train_processed['buy_qty'] / \n                                                (train_processed['sell_qty'] + 1)).astype(np.float32)\n            test_processed['buy_sell_ratio'] = (test_processed['buy_qty'] / \n                                               (test_processed['sell_qty'] + 1)).astype(np.float32)\n            \n            # Handle infinities\n            train_processed['buy_sell_ratio'] = train_processed['buy_sell_ratio'].replace([np.inf, -np.inf], 0)\n            test_processed['buy_sell_ratio'] = test_processed['buy_sell_ratio'].replace([np.inf, -np.inf], 0)\n            \n            # Log transformations\n            train_processed['log_buy_qty'] = np.log1p(train_processed['buy_qty']).astype(np.float32)\n            test_processed['log_buy_qty'] = np.log1p(test_processed['buy_qty']).astype(np.float32)\n            \n            train_processed['log_sell_qty'] = np.log1p(train_processed['sell_qty']).astype(np.float32)\n            test_processed['log_sell_qty'] = np.log1p(test_processed['sell_qty']).astype(np.float32)\n        \n        # 3. Volume features\n        if 'volume' in train.columns:\n            # Fill missing values\n            train_processed['volume'] = train_processed['volume'].fillna(0)\n            test_processed['volume'] = test_processed['volume'].fillna(0)\n            \n            # Log volume\n            train_processed['log_volume'] = np.log1p(train_processed['volume']).astype(np.float32)\n            test_processed['log_volume'] = np.log1p(test_processed['volume']).astype(np.float32)\n            \n            # Sqrt volume (dampens extreme values)\n            train_processed['sqrt_volume'] = np.sqrt(train_processed['volume']).astype(np.float32)\n            test_processed['sqrt_volume'] = np.sqrt(test_processed['volume']).astype(np.float32)\n            \n            # Volume vs total order book\n            if 'bid_qty' in train.columns and 'ask_qty' in train.columns:\n                train_processed['volume_to_orderbook'] = (train_processed['volume'] / \n                                                         (train_processed['bid_qty'] + train_processed['ask_qty'] + 1)).astype(np.float32)\n                test_processed['volume_to_orderbook'] = (test_processed['volume'] / \n                                                        (test_processed['bid_qty'] + test_processed['ask_qty'] + 1)).astype(np.float32)\n                \n                # Handle infinities\n                train_processed['volume_to_orderbook'] = train_processed['volume_to_orderbook'].replace([np.inf, -np.inf], 0)\n                test_processed['volume_to_orderbook'] = test_processed['volume_to_orderbook'].replace([np.inf, -np.inf], 0)\n        \n        # 4. Liquidity features\n        if all(col in train.columns for col in ['bid_qty', 'ask_qty', 'volume']):\n            # Total liquidity\n            train_processed['total_liquidity'] = (train_processed['bid_qty'] + train_processed['ask_qty']).astype(np.float32)\n            test_processed['total_liquidity'] = (test_processed['bid_qty'] + test_processed['ask_qty']).astype(np.float32)\n            \n            # Log liquidity\n            train_processed['log_liquidity'] = np.log1p(train_processed['total_liquidity']).astype(np.float32)\n            test_processed['log_liquidity'] = np.log1p(test_processed['total_liquidity']).astype(np.float32)\n            \n            # Liquidity consumption (volume as fraction of liquidity)\n            train_processed['liquidity_consumption'] = (train_processed['volume'] / \n                                                      (train_processed['total_liquidity'] + 1)).astype(np.float32)\n            test_processed['liquidity_consumption'] = (test_processed['volume'] / \n                                                     (test_processed['total_liquidity'] + 1)).astype(np.float32)\n            \n            # Handle infinities\n            train_processed['liquidity_consumption'] = train_processed['liquidity_consumption'].replace([np.inf, -np.inf], 0)\n            test_processed['liquidity_consumption'] = test_processed['liquidity_consumption'].replace([np.inf, -np.inf], 0)\n        \n        # 5. Pressure indicators\n        if all(col in train.columns for col in ['buy_qty', 'sell_qty', 'bid_qty', 'ask_qty']):\n            # Buy pressure\n            train_processed['buy_pressure'] = ((train_processed['buy_qty'] + train_processed['bid_qty']) / \n                                             (train_processed['buy_qty'] + train_processed['bid_qty'] + \n                                              train_processed['sell_qty'] + train_processed['ask_qty'] + 1)).astype(np.float32)\n            test_processed['buy_pressure'] = ((test_processed['buy_qty'] + test_processed['bid_qty']) / \n                                            (test_processed['buy_qty'] + test_processed['bid_qty'] + \n                                             test_processed['sell_qty'] + test_processed['ask_qty'] + 1)).astype(np.float32)\n            \n            # Net order flow\n            train_processed['net_order_flow'] = (train_processed['buy_qty'] - train_processed['sell_qty'] + \n                                                train_processed['bid_qty'] - train_processed['ask_qty']).astype(np.float32)\n            test_processed['net_order_flow'] = (test_processed['buy_qty'] - test_processed['sell_qty'] + \n                                               test_processed['bid_qty'] - test_processed['ask_qty']).astype(np.float32)\n        \n        # 6. Statistical transformations of anonymous features\n        # Select a subset of anonymous features for advanced transformations\n        selected_anon_features = [f for f in self.anon_features if f in train.columns][:50]  # Use first 50\n        \n        if selected_anon_features:\n            # Create statistical aggregates\n            train_processed['anon_mean'] = train_processed[selected_anon_features].mean(axis=1).astype(np.float32)\n            test_processed['anon_mean'] = test_processed[selected_anon_features].mean(axis=1).astype(np.float32)\n            \n            train_processed['anon_std'] = train_processed[selected_anon_features].std(axis=1).astype(np.float32)\n            test_processed['anon_std'] = test_processed[selected_anon_features].std(axis=1).astype(np.float32)\n            \n            train_processed['anon_max'] = train_processed[selected_anon_features].max(axis=1).astype(np.float32)\n            test_processed['anon_max'] = test_processed[selected_anon_features].max(axis=1).astype(np.float32)\n            \n            train_processed['anon_min'] = train_processed[selected_anon_features].min(axis=1).astype(np.float32)\n            test_processed['anon_min'] = test_processed[selected_anon_features].min(axis=1).astype(np.float32)\n            \n            train_processed['anon_range'] = (train_processed['anon_max'] - train_processed['anon_min']).astype(np.float32)\n            test_processed['anon_range'] = (test_processed['anon_max'] - test_processed['anon_min']).astype(np.float32)\n        \n        # 7. Power transformations for key features\n        power_features = ['volume', 'bid_qty', 'ask_qty', 'buy_qty', 'sell_qty']\n        power_features = [f for f in power_features if f in train.columns]\n        \n        if power_features:\n            pt = PowerTransformer(method='yeo-johnson')\n            \n            for feature in power_features:\n                # Skip if all values are the same\n                if train_processed[feature].std() > 0:\n                    train_values = train_processed[feature].values.reshape(-1, 1)\n                    test_values = test_processed[feature].values.reshape(-1, 1)\n                    \n                    train_transformed = pt.fit_transform(train_values)\n                    test_transformed = pt.transform(test_values)\n                    \n                    train_processed[f'{feature}_yj'] = train_transformed.ravel().astype(np.float32)\n                    test_processed[f'{feature}_yj'] = test_transformed.ravel().astype(np.float32)\n        \n        # 8. Polynomial features for market microstructure\n        if len(self.market_features) >= 2:\n            poly_features = []\n            for feature in self.market_features:\n                if feature in train.columns:\n                    poly_features.append(feature)\n            \n            if len(poly_features) >= 2:\n                # Create polynomial features (degree 2 for interactions)\n                poly = PolynomialFeatures(degree=2, include_bias=False, interaction_only=True)\n                \n                # Fill missing values before polynomial transformation\n                train_poly_data = train_processed[poly_features].fillna(0)\n                test_poly_data = test_processed[poly_features].fillna(0)\n                \n                # Fit and transform\n                poly_train = poly.fit_transform(train_poly_data)\n                poly_test = poly.transform(test_poly_data)\n                \n                # Add polynomial features\n                n_orig = len(poly_features)\n                for i in range(n_orig, poly_train.shape[1]):\n                    train_processed[f'poly_{i}'] = poly_train[:, i].astype(np.float32)\n                    test_processed[f'poly_{i}'] = poly_test[:, i].astype(np.float32)\n        \n        # 9. K-means clustering on market features\n        cluster_features = ['bid_qty', 'ask_qty', 'buy_qty', 'sell_qty', 'volume']\n        cluster_features = [f for f in cluster_features if f in train.columns]\n        \n        if len(cluster_features) >= 2:\n            # Extract cluster data and handle missing values\n            cluster_train = train_processed[cluster_features].fillna(0)\n            cluster_test = test_processed[cluster_features].fillna(0)\n            \n            # Normalize data for clustering\n            scaler = RobustScaler()\n            cluster_train_scaled = scaler.fit_transform(cluster_train)\n            cluster_test_scaled = scaler.transform(cluster_test)\n            \n            # Apply K-means clustering\n            kmeans = KMeans(n_clusters=5, random_state=self.seed, n_init=10)\n            train_processed['market_cluster'] = kmeans.fit_predict(cluster_train_scaled).astype(np.uint8)\n            test_processed['market_cluster'] = kmeans.predict(cluster_test_scaled).astype(np.uint8)\n            \n            # Create dummy variables for clusters\n            for i in range(5):\n                train_processed[f'cluster_{i}'] = (train_processed['market_cluster'] == i).astype(np.uint8)\n                test_processed[f'cluster_{i}'] = (test_processed['market_cluster'] == i).astype(np.uint8)\n            \n            # Cluster distances\n            cluster_distances = kmeans.transform(cluster_train_scaled)\n            cluster_distances_test = kmeans.transform(cluster_test_scaled)\n            \n            train_processed['min_cluster_dist'] = np.min(cluster_distances, axis=1).astype(np.float32)\n            test_processed['min_cluster_dist'] = np.min(cluster_distances_test, axis=1).astype(np.float32)\n        \n        # Prepare for modeling - remove ID and target columns\n        train_processed = train_processed.drop([self.id_column, self.target_column], axis=1)\n        test_processed = test_processed.drop([self.id_column], axis=1)\n        \n        # Log feature engineering statistics\n        print(f\"{get_timestamp()} - INFO - Feature engineering complete - created {train_processed.shape[1]} features\")\n        \n        # Clean up memory\n        gc.collect()\n        \n        return train_processed, test_processed\n    \n    def _analyze_target_distribution(self, target: pd.Series) -> None:\n        \"\"\"\n        Analyze target distribution to identify outlier patterns.\n        \n        Args:\n            target: Series containing target values\n        \"\"\"\n        # Basic statistics\n        print(f\"{get_timestamp()} - INFO - Target distribution statistics:\")\n        print(f\"{get_timestamp()} - INFO -   Min: {target.min():.4f}\")\n        print(f\"{get_timestamp()} - INFO -   Max: {target.max():.4f}\")\n        print(f\"{get_timestamp()} - INFO -   Mean: {target.mean():.4f}\")\n        print(f\"{get_timestamp()} - INFO -   Median: {target.median():.4f}\")\n        print(f\"{get_timestamp()} - INFO -   Std Dev: {target.std():.4f}\")\n        \n        # Percentile analysis\n        percentiles = [1, 5, 25, 50, 75, 95, 99, 99.5, 99.9]\n        for p in percentiles:\n            value = np.percentile(target, p)\n            print(f\"{get_timestamp()} - INFO -   {p}th percentile: {value:.4f}\")\n    \n    def preprocess_data(self, train_processed: pd.DataFrame, test_processed: pd.DataFrame) -> bool:\n        \"\"\"\n        Preprocess data for modeling with optimized scaling and feature selection.\n        \n        Args:\n            train_processed: Engineered training DataFrame\n            test_processed: Engineered test DataFrame\n            \n        Returns:\n            Boolean indicating successful preprocessing\n        \"\"\"\n        print(f\"{get_timestamp()} - INFO - Preprocessing data with optimization...\")\n        \n        # All columns should be numeric at this point\n        numerical_cols = train_processed.columns.tolist()\n        print(f\"{get_timestamp()} - INFO - Processing {len(numerical_cols)} numerical features\")\n        \n        # Apply Robust scaling for all numerical features (handles outliers)\n        if len(numerical_cols) > 0:\n            scaler = RobustScaler()\n            \n            # Scale\n            train_numerical = train_processed[numerical_cols].fillna(0)\n            test_numerical = test_processed[numerical_cols].fillna(0)\n            \n            # Fit and transform\n            train_processed[numerical_cols] = scaler.fit_transform(train_numerical).astype(np.float32)\n            test_processed[numerical_cols] = scaler.transform(test_numerical).astype(np.float32)\n        \n        # Handle any remaining missing values\n        train_processed = train_processed.fillna(0)\n        test_processed = test_processed.fillna(0)\n        \n        # Remove low variance features\n        if train_processed.shape[1] > 10:\n            print(f\"{get_timestamp()} - INFO - Removing low variance features for model efficiency...\")\n            \n            selector = VarianceThreshold(threshold=0.02)\n            \n            # Fit and transform\n            train_processed_filtered = selector.fit_transform(train_processed)\n            test_processed_filtered = selector.transform(test_processed)\n            \n            # Get feature names after selection\n            selected_features = train_processed.columns[selector.get_support()]\n            \n            # Convert back to DataFrame with column names\n            train_processed = pd.DataFrame(train_processed_filtered, columns=selected_features)\n            test_processed = pd.DataFrame(test_processed_filtered, columns=selected_features)\n            \n            # Log feature reduction results\n            print(f\"{get_timestamp()} - INFO - Kept {len(selected_features)} features after variance filtering\")\n        \n        # Remove highly correlated features\n        if train_processed.shape[1] > 20:\n            print(f\"{get_timestamp()} - INFO - Removing highly correlated features...\")\n            \n            # Calculate correlation matrix\n            corr_matrix = train_processed.corr().abs()\n            \n            # Create upper triangle mask\n            upper = corr_matrix.where(np.triu(np.ones(corr_matrix.shape), k=1).astype(bool))\n            \n            # Find features with correlation greater than threshold\n            to_drop = [column for column in upper.columns if any(upper[column] > 0.95)]\n            \n            if to_drop:\n                # Drop highly correlated features\n                train_processed = train_processed.drop(to_drop, axis=1)\n                test_processed = test_processed.drop(to_drop, axis=1)\n                \n                print(f\"{get_timestamp()} - INFO - Removed {len(to_drop)} highly correlated features\")\n        \n        # Store processed data\n        self.X_train = train_processed\n        self.X_test = test_processed\n        \n        # Clean up memory \n        gc.collect()\n        \n        print(f\"{get_timestamp()} - INFO - Preprocessing complete with {train_processed.shape[1]} final features\")\n        return True\n    \n    def initialize_ray(self) -> bool:\n        \"\"\"\n        Initialize Ray for distributed computing.\n        \n        Returns:\n            Boolean indicating successful initialization\n        \"\"\"\n        print(f\"{get_timestamp()} - INFO - Initializing Ray with {self.num_cpus} CPUs...\")\n        \n        if not ray.is_initialized():\n            # Configure Ray with specified CPUs\n            ray.init(\n                num_cpus=self.num_cpus, \n                ignore_reinit_error=True, \n                log_to_driver=False\n            )\n            \n            # Log Ray cluster resources\n            resources = ray.cluster_resources()\n            print(f\"{get_timestamp()} - INFO - Ray initialized with resources: {resources}\")\n        else:\n            print(f\"{get_timestamp()} - INFO - Ray was already initialized.\")\n            \n        return True\n    \n    def shutdown_ray(self) -> None:\n        \"\"\"\n        Properly shutdown Ray to release resources.\n        \"\"\"\n        if ray.is_initialized():\n            print(f\"{get_timestamp()} - INFO - Shutting down Ray...\")\n            ray.shutdown()\n            print(f\"{get_timestamp()} - INFO - Ray resources released\")\n    \n    def train_model(self) -> bool:\n        \"\"\"\n        Train model using FLAML AutoML with Ray for distributed hyperparameter tuning.\n        \n        Returns:\n            Boolean indicating successful training\n        \"\"\"\n        print(f\"{get_timestamp()} - INFO - Training FLAML model with enhanced configuration...\")\n        \n        # Initialize Ray\n        ray_initialized = self.initialize_ray()\n        if not ray_initialized:\n            print(f\"{get_timestamp()} - ERROR - Failed to initialize Ray. Falling back to non-distributed training.\")\n        \n        # Set random seed for reproducibility\n        np.random.seed(self.seed)\n        \n        # Create and configure FLAML\n        automl = AutoML()\n        \n        # Configure distributed training with Ray and enhanced settings\n        settings = {\n            \"time_budget\": self.flaml_time_budget,\n            \"task\": \"regression\",\n            \"metric\": \"mse\",  # Use MSE for log-transformed target\n            \"n_concurrent_trials\": self.num_cpus,\n            \"free_mem_ratio\": 0.1,\n            \"sample\": False\n        }\n        \n        # Train the model\n        start_training = time.time()\n        \n        automl.fit(\n            X_train=self.X_train,\n            y_train=self.y_train,\n            **settings\n        )\n        \n        training_time = (time.time() - start_training) / 60\n        \n        # Store the model\n        self.model = automl\n        \n        # Log model information\n        print(f\"{get_timestamp()} - INFO - Training complete in {training_time:.2f} minutes\")\n        print(f\"{get_timestamp()} - INFO - Best model: {automl.best_estimator}\")\n        print(f\"{get_timestamp()} - INFO - Best configuration: {automl.best_config}\")\n        print(f\"{get_timestamp()} - INFO - Best validation score (MSE): {automl.best_loss:.5f}\")\n        print(f\"{get_timestamp()} - INFO - Training time for best model: {automl.best_config_train_time:.2f} seconds\")\n        \n        # Print feature importance if available\n        if hasattr(automl.model, 'feature_importances_'):\n            importances = automl.model.feature_importances_\n            feature_names = self.X_train.columns\n            importance_df = pd.DataFrame({'feature': feature_names, 'importance': importances})\n            importance_df = importance_df.sort_values('importance', ascending=False)\n            print(f\"{get_timestamp()} - INFO - Top 10 important features:\")\n            for i, row in importance_df.head(10).iterrows():\n                print(f\"{get_timestamp()} - INFO -   {row['feature']}: {row['importance']:.4f}\")\n        \n        # Ensure Ray is shut down properly to release resources\n        self.shutdown_ray()\n        \n        return True\n    \n    def generate_predictions(self) -> bool:\n        \"\"\"\n        Generate predictions and create submission file.\n        \n        Returns:\n            Boolean indicating successful prediction generation\n        \"\"\"\n        print(f\"{get_timestamp()} - INFO - Generating predictions...\")\n        \n        # Get predictions\n        predictions = self.model.predict(self.X_test)\n        \n        # Create submission file\n        submission = pd.DataFrame({\n            self.id_column: self.test_ids,\n            self.target_column: predictions\n        })\n        \n        # Save to CSV\n        submission_path = os.path.join(self.output_path, 'submission.csv')\n        submission.to_csv(submission_path, index=False)\n        \n        print(f\"{get_timestamp()} - INFO - Submission file created at {submission_path}\")\n        \n        # Generate basic prediction statistics\n        print(f\"{get_timestamp()} - INFO - Prediction statistics:\")\n        print(f\"{get_timestamp()} - INFO -   Count: {len(predictions)}\")\n        print(f\"{get_timestamp()} - INFO -   Min: {predictions.min():.4f}\")\n        print(f\"{get_timestamp()} - INFO -   Max: {predictions.max():.4f}\")\n        print(f\"{get_timestamp()} - INFO -   Mean: {predictions.mean():.4f}\")\n        print(f\"{get_timestamp()} - INFO -   Median: {np.median(predictions):.4f}\")\n        print(f\"{get_timestamp()} - INFO -   Std Dev: {np.std(predictions):.4f}\")\n        \n        # Compare with original training data distribution\n        print(f\"{get_timestamp()} - INFO - Training target statistics for comparison:\")\n        print(f\"{get_timestamp()} - INFO -   Min: {self.y_train_original.min():.4f}\")\n        print(f\"{get_timestamp()} - INFO -   Max: {self.y_train_original.max():.4f}\")\n        print(f\"{get_timestamp()} - INFO -   Mean: {self.y_train_original.mean():.4f}\")\n        print(f\"{get_timestamp()} - INFO -   Median: {self.y_train_original.median():.4f}\")\n        print(f\"{get_timestamp()} - INFO -   Std Dev: {self.y_train_original.std():.4f}\")\n        \n        return True\n    \n    def run(self) -> bool:\n        \"\"\"\n        Execute the complete prediction pipeline with robust error handling.\n        \n        Returns:\n            Boolean indicating successful pipeline execution\n        \"\"\"\n        # Record start time for overall process\n        start_time = time.time()\n        \n        # Set random seed for reproducibility\n        np.random.seed(self.seed)\n        \n        print(f\"{get_timestamp()} - INFO - Starting crypto price movement prediction pipeline\")\n        \n        # Step 1: Load data\n        train, test = self.load_data()\n        \n        # Step 2: Feature engineering\n        train_processed, test_processed = self.engineer_features(train, test)\n        \n        # Step 3: Preprocess data\n        preprocess_success = self.preprocess_data(train_processed, test_processed)\n        if not preprocess_success:\n            print(f\"{get_timestamp()} - ERROR - Data preprocessing failed. Exiting.\")\n            return False\n        \n        # Step 4: Train model with Ray\n        train_success = self.train_model()\n        if not train_success:\n            print(f\"{get_timestamp()} - ERROR - Model training failed. Exiting.\")\n            return False\n        \n        # Step 5: Generate predictions\n        predict_success = self.generate_predictions()\n        if not predict_success:\n            print(f\"{get_timestamp()} - ERROR - Prediction generation failed. Exiting.\")\n            return False\n        \n        # Calculate and log total runtime\n        total_time = (time.time() - start_time) / 60\n        print(f\"{get_timestamp()} - INFO - Pipeline completed successfully in {total_time:.2f} minutes\")\n        \n        # Ensure Ray is shut down properly\n        self.shutdown_ray()\n        \n        return True\n\n\n# Run the pipeline\nif __name__ == \"__main__\":\n    # Adjust these paths for your environment\n    input_path = '/kaggle/input/drw-quantitative-research-challenge/'\n    output_path = '/kaggle/working/'\n    \n    # Create predictor with CPUs for distributed computing\n    predictor = CryptoPricePredictor(\n        input_path=input_path, \n        output_path=output_path, \n        seed=42,\n        num_cpus=4  # Adjust based on available computational resources\n    )\n    \n    # Run the pipeline\n    predictor.run()","metadata":{"trusted":true},"outputs":[],"execution_count":null}]}