{"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":12993472,"sourceType":"competition"},{"sourceId":12405352,"sourceType":"datasetVersion","datasetId":7823194},{"sourceId":249869065,"sourceType":"kernelVersion"}],"dockerImageVersionId":31040,"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":"#!/usr/bin/env python\n# coding: utf-8\n\n# 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\nimport numpy as np # linear algebra\nimport pandas as pd # data processing, CSV file I/O (e.g. pd.read_csv)\nimport matplotlib.pyplot as plt\nfrom sklearn.linear_model import Ridge, SGDRegressor\nfrom sklearn.preprocessing import StandardScaler\nfrom sklearn.ensemble import VotingRegressor\nimport os\nimport gc\nimport warnings\nfrom collections import defaultdict\nwarnings.filterwarnings('ignore')\n\n# Input data files are available in the read-only \"../input/\" directory\nfor dirname, _, filenames in os.walk('/kaggle/input'):\n    for filename in filenames:\n        print(os.path.join(dirname, filename))\n\ndef optimize_memory(df, verbose=True):\n    \"\"\"\n    Optimize memory usage by downcasting numeric types where possible.\n    \"\"\"\n    if verbose:\n        start_mem = df.memory_usage().sum() / 1024**2\n        print(f'Memory usage before optimization: {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.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 optimization: {end_mem:.2f} MB')\n        print(f'Decreased by {100 * (start_mem - end_mem) / start_mem:.1f}%')\n    \n    return df\n\ndef get_feature_columns(df):\n    \"\"\"Get all feature columns including the extended list and new variables.\"\"\"\n    # Extended feature list provided\n    extended_features = [\n        'X727', 'X427', 'X288', 'X721', 'X312', 'X421', 'X471', 'X573', 'X780', 'X255',\n        'X144', 'X299', 'X301', 'X563', 'X737', 'X702', 'ask_qty', 'X507', 'X306', 'X501',\n        'X303', 'amihud_illiquidity', 'X586', 'X43', 'X517', 'X248', 'X137', 'X757', 'X196',\n        'X777', 'X280', 'X266', 'X689', 'X294', 'X492', 'X555', 'X731', 'X262', 'X576',\n        'X13', 'X518', 'X502', 'X558', 'pin_proxy', 'X6', 'X602', 'X695', 'X703', 'X413',\n        'X660', 'X37', 'X15', 'X310', 'X512', 'X362', 'X631', 'X214', 'X562', 'X488',\n        'X510', 'X256', 'X35', 'X128', 'X86', 'X170', 'X30', 'X265', 'X323', 'X559',\n        'X348', 'X130', 'X529', 'X20', 'X4', 'X90', 'X192', 'X91', 'X582', 'X99',\n        'X24', 'X317', 'X707', 'X653', 'X519', 'X557', 'X371', 'X415', 'X84', 'X83',\n        'order_toxicity', 'X360', 'X111', 'X699', 'X187', 'X591', 'X637', 'X567', 'X577',\n        'X313', 'X60', 'X671', 'X698', 'X701', 'X725', 'X292', 'X638', 'X741', 'X379',\n        'X700', 'X614', 'X676', 'X516', 'X697', 'X611', 'X311', 'X615', 'X706', 'X466',\n        'X571', 'X451', 'X17', 'X584', 'X436', 'X305', 'liquidity_consumption', 'X34', 'X282',\n        'X681', 'X7', 'X208', 'X41', 'X536', 'X548', 'X296', 'X776', 'X87', 'X40',\n        'X570', 'X539', 'X474', 'X753', 'X425', 'X217', 'X199', 'X18', 'X609', 'X21',\n        'X277', 'X279', 'X326', 'X540', 'X688', 'X553', 'X452', 'X738', 'X183', 'X759',\n        'bid_ask_ratio', 'X495', 'volume_participation', 'X715', 'X385', 'X291', 'X409', 'X112',\n        'X693', 'X102', 'X318', 'X705', 'X556', 'X547'\n    ]\n    \n    # Core features that are usually important\n    core_features = [\n        'X363', 'X405', 'X321', 'X175', 'X179', 'X197', 'X22', 'X181', 'X28', 'X169',\n        'X198', 'X173', 'X338', 'X344', 'X587', 'X450', 'X97', 'X52', 'X444', 'X598',\n        'X297', 'X138', 'X572', 'X343', 'X438', 'X459', 'X758', 'X25', 'buy_qty', \n        'sell_qty', 'volume', 'bid_qty'\n    ]\n    \n    # NEW features to add\n    new_features = [\n        'X363', 'X321', 'X405', 'X730', 'X523', 'X756', 'X589', 'X462', 'X779', 'log_liquidity',\n        'X25', 'X532', 'X520', 'X329', 'X383', 'X751', 'X535', 'X639', 'X596', 'X761',\n        'X145', 'X709', 'X173', 'X245', 'X168', 'X171', 'X241', 'X31', 'X105', 'X63',\n        'X263', 'X426', 'X286', 'X357', 'X399', 'X315', 'X468', 'X131', 'X647', 'log_spread',\n        'X752', 'X254', 'X592', 'X733', 'X636', 'X394', 'X527', 'X180', 'X367', 'X38',\n        'X634', 'X718', 'X387', 'X429', 'X345', 'X344', 'X253', 'X469', 'X446', 'X125',\n        'X760', 'X186', 'X711', 'X150', 'X661', 'X215', 'X403', 'X141', 'X771', 'X453',\n        'X401', 'X629', 'X616', 'X281', 'X432', 'X283', 'X244', 'X440', 'X430', 'X382',\n        'X175', 'X95', 'X444', 'X189', 'X55', 'X605', 'X663', 'X194', 'X439', 'X670',\n        'X483', 'X163', 'X376', 'X71', 'X650', 'X203', 'X8', 'X624', 'X160', 'X100',\n        'X14', 'X511', 'X59', 'X302', 'X81', 'X325', 'X514', 'X649', 'X447', 'X538',\n        'X443', 'X39', 'X343', 'X12', 'X678', 'X775', 'X498', 'X249', 'X42', 'X384',\n        'kyle_lambda', 'X349', 'X356', 'X2', 'X250', 'X397', 'X685', 'X568', 'X136', 'X496',\n        'X53', 'X66', 'X374', 'X590', 'X668', 'X585', 'X677', 'X667', 'X530', 'X28',\n        'X64', 'X407', 'X494', 'X770', 'X710', 'X526', 'X644', 'X167', 'X190', 'X723',\n        'X33', 'X579', 'X206'\n    ]\n    \n    # ADDITIONAL features requested by user\n    additional_features = [\n        'X525', 'X267', 'X166', 'X719', 'X489', 'X758', 'X652', 'X433', 'X778', 'X428',\n        'X617', 'X259', 'X633', 'X565', 'X364', 'depth_ratio', 'X550', 'X687', 'X610', 'X599',\n        'X717', 'X587', 'X143', 'X506', 'X546', 'X505', 'X159', 'X574', 'X278', 'X458',\n        'X1', 'X749', 'X155', 'X651', 'X470', 'X580', 'X445', 'X373', 'X82', 'X607',\n        'X298', 'X221', 'X388', 'X120', 'X391', 'X23', 'X679', 'X377', 'X767', 'X755',\n        'X566', 'X424', 'X438', 'X198', 'X300', 'X268', 'X434', 'X290', 'X368', 'X464',\n        'X119', 'X197', 'X597', 'X157', 'X485', 'X127', 'X101', 'X533', 'X235', 'X712',\n        'X154', 'X239', 'X10', 'X420', 'X449', 'X740', 'X227', 'X36', 'X358', 'X551',\n        'X528', 'X285', 'X335', 'X152', 'X110', 'X68', 'X713', 'X402', 'X370', 'X735',\n        'X200', 'X331', 'X473', 'X162', 'X213', 'X322', 'X289', 'X477', 'X113', 'X560',\n        'X672', 'X621', 'X682', 'X5', 'X72', 'X44', 'X419', 'buy_pressure', 'X242', 'volume',\n        'X472', 'X332', 'X441', 'buy_sell_ratio', 'pressure_ratio', 'X508', 'X594', 'X191',\n        'X261', 'X603', 'net_pressure', 'order_flow_imbalance', 'sell_pressure', 'X240', 'X673',\n        'X608', 'X509', 'X165', 'X720', 'X314', 'X522', 'X531', 'X625', 'bid_depth_ratio',\n        'X435', 'X293', 'X486', 'price_efficiency', 'X716', 'X627', 'X626', 'X169', 'X613',\n        'X680', 'X544', 'X115', 'X307', 'X665', 'X465', 'X347', 'X728', 'X70', 'log_volume',\n        'X340', 'X459', 'X56', 'X395', 'X354', 'X51', 'X732', 'X247', 'X324', 'X316',\n        'X76', 'X341', 'X739', 'X601', 'X386', 'X683', 'X149', 'X193', 'X628', 'X309',\n        'X351', 'X393'\n    ]\n    \n    # NEW: 100 Additional X features\n    additional_100_features = [\n        # Features from X1-X100 range not yet included\n        'X3', 'X9', 'X11', 'X16', 'X19', 'X26', 'X27', 'X29', 'X32', 'X45',\n        'X46', 'X47', 'X48', 'X49', 'X50', 'X54', 'X57', 'X58', 'X61', 'X62',\n        'X65', 'X67', 'X69', 'X73', 'X74', 'X75', 'X77', 'X78', 'X79', 'X80',\n        'X85', 'X88', 'X89', 'X92', 'X93', 'X94', 'X96', 'X98',\n        \n        # Features from X100-X300 range\n        'X103', 'X104', 'X106', 'X107', 'X108', 'X109', 'X114', 'X116', 'X117', 'X118',\n        'X121', 'X122', 'X123', 'X124', 'X126', 'X129', 'X132', 'X133', 'X134', 'X135',\n        'X139', 'X140', 'X142', 'X146', 'X147', 'X148', 'X151', 'X153', 'X156', 'X158',\n        'X161', 'X164', 'X172', 'X174', 'X176', 'X177', 'X178', 'X182', 'X184', 'X185',\n        'X188', 'X195', 'X201', 'X202', 'X204', 'X205', 'X207', 'X209', 'X210', 'X211',\n        'X212', 'X216', 'X218', 'X219', 'X220', 'X222', 'X223', 'X224', 'X225', 'X226',\n        'X228', 'X229', 'X230', 'X231', 'X232', 'X233', 'X234', 'X236', 'X237', 'X238',\n        'X243', 'X246', 'X251', 'X252', 'X257', 'X258', 'X260', 'X264', 'X269', 'X270',\n        'X271', 'X272', 'X273', 'X275', 'X276', 'X284', 'X287', 'X295', 'X304', 'X308',\n        'X319', 'X320', 'X327', 'X328', 'X330', 'X333', 'X334', 'X336', 'X337', 'X339'\n    ]\n    \n    # Combine all features and remove duplicates\n    all_features = list(dict.fromkeys(\n        extended_features + core_features + new_features + \n        additional_features + additional_100_features\n    ))\n    \n    # Filter to only available columns\n    available_features = [col for col in all_features if col in df.columns]\n    \n    print(f\"Found {len(available_features)} features out of {len(all_features)} requested\")\n    \n    return available_features\n\nclass IncrementalLagEnsemble:\n    \"\"\"Incrementally trains models on different lag configurations using SGD.\"\"\"\n    \n    def __init__(self, feature_batch_size=50, lag_batch_size=5, n_epochs=9):\n        self.feature_batch_size = feature_batch_size\n        self.lag_batch_size = lag_batch_size\n        self.n_epochs = n_epochs  # Increased to 9 as requested\n        self.models = {}\n        self.scalers = {}\n        self.feature_names = None\n        self.model_weights = {}\n        self.performance_history = defaultdict(list)\n        \n        # Define lag strategies with odd/even splits for many configurations\n        self.lag_strategies = {\n            # Original micro lags\n            'micro': [1, 2, 3, 4, 5],\n            'micro_odd': [1, 3, 5, 7, 9],\n            'micro_even': [2, 4, 6, 8, 10],\n            \n            # Ultra short with odd/even\n            'ultra_short': [6, 8, 10, 12, 15],\n            'ultra_short_odd': [7, 9, 11, 13, 15, 17],\n            'ultra_short_even': [6, 8, 10, 12, 14, 16],\n            \n            # Short with odd/even\n            'short': [20, 25, 30, 40, 50],\n            'short_odd': [21, 25, 31, 41, 51],\n            'short_even': [20, 24, 30, 40, 50],\n            \n            # Short medium with odd/even\n            'short_medium': [60, 75, 90, 105, 120],\n            'short_medium_odd': [61, 75, 91, 105, 121],\n            'short_medium_even': [60, 74, 90, 104, 120],\n            \n            # Medium with odd/even\n            'medium': [150, 180, 210, 240, 300],\n            'medium_odd': [151, 181, 211, 241, 301],\n            'medium_even': [150, 180, 210, 240, 300],\n            \n            # Medium long with odd/even\n            'medium_long': [360, 420, 480, 540, 600],\n            'medium_long_odd': [361, 421, 481, 541, 601],\n            'medium_long_even': [360, 420, 480, 540, 600],\n            \n            # Long with odd/even\n            'long': [720, 840, 960, 1080, 1200],\n            'long_odd': [721, 841, 961, 1081, 1201],\n            'long_even': [720, 840, 960, 1080, 1200],\n            \n            # Very long with odd/even\n            'very_long': [1440, 1800, 2160, 2520, 2880],\n            'very_long_odd': [1441, 1801, 2161, 2521, 2881],\n            'very_long_even': [1440, 1800, 2160, 2520, 2880],\n            \n            # Ultra long (keeping original only due to very large values)\n            'ultra_long': [3600, 4320, 5040, 5760, 7200]\n        }\n        \n    def create_lag_features_batch(self, df, lag_list):\n        \"\"\"Create lag features for a batch of lags.\"\"\"\n        lag_features = []\n        \n        for lag in lag_list:\n            lagged = df.shift(-lag)\n            lagged.columns = [f'{col}_lag_{lag}' for col in df.columns]\n            lag_features.append(lagged)\n        \n        result = pd.concat([df] + lag_features, axis=1)\n        result = result.fillna(0)\n        \n        return result\n    \n    def train_sgd_model(self, X, y, model_name):\n        \"\"\"Train SGD model incrementally with more epochs.\"\"\"\n        print(f\"  Training SGD model: {model_name} ({self.n_epochs} epochs)\")\n        \n        # Initialize model and scaler if not exists\n        if model_name not in self.models:\n            self.models[model_name] = SGDRegressor(\n                loss='huber',\n                penalty='elasticnet',\n                alpha=0.0001,\n                l1_ratio=0.15,\n                learning_rate='invscaling',\n                eta0=0.01,\n                power_t=0.25,\n                random_state=42,\n                warm_start=True,\n                max_iter=1000,\n                tol=1e-3\n            )\n            self.scalers[model_name] = StandardScaler()\n            \n        model = self.models[model_name]\n        scaler = self.scalers[model_name]\n        \n        # Train in epochs with smaller chunks for better convergence\n        chunk_size = 25000  # Slightly smaller chunks for more updates\n        for epoch in range(self.n_epochs):\n            # Shuffle indices for each epoch\n            indices = np.random.permutation(len(X))\n            \n            for start_idx in range(0, len(X), chunk_size):\n                end_idx = min(start_idx + chunk_size, len(X))\n                \n                # Get shuffled chunk\n                chunk_indices = indices[start_idx:end_idx]\n                X_chunk = X.iloc[chunk_indices]\n                y_chunk = y[chunk_indices]\n                \n                # Scale\n                if start_idx == 0 and epoch == 0:\n                    X_scaled = scaler.fit_transform(X_chunk)\n                else:\n                    X_scaled = scaler.transform(X_chunk)\n                \n                # Partial fit\n                model.partial_fit(X_scaled, y_chunk)\n            \n            # Print progress\n            if epoch % 3 == 0:\n                print(f\"    Epoch {epoch+1}/{self.n_epochs} completed\")\n        \n        return model\n    \n    def train_feature_batch(self, feature_batch, X_full, y, strategy_name, lag_list):\n        \"\"\"Train on a batch of features with specific lags.\"\"\"\n        print(f\"\\n  Processing feature batch ({len(feature_batch)} features) with {strategy_name} lags\")\n        \n        # Select feature batch\n        X_batch = X_full[feature_batch].copy()\n        \n        # Create lag features\n        X_with_lags = self.create_lag_features_batch(X_batch, lag_list)\n        \n        # Train SGD model\n        model_name = f\"{strategy_name}_{feature_batch[0]}_{feature_batch[-1]}\"\n        self.train_sgd_model(X_with_lags, y, model_name)\n        \n        # Clean up\n        del X_batch, X_with_lags\n        gc.collect()\n        \n    def fit(self, X, y):\n        \"\"\"Fit ensemble using incremental training.\"\"\"\n        print(\"Training Incremental Lag Ensemble with Odd/Even Lag Configurations...\")\n        print(f\"Total epochs per model: {self.n_epochs}\")\n        \n        self.feature_names = X.columns.tolist()\n        n_features = len(self.feature_names)\n        \n        # Split features into batches\n        feature_batches = []\n        for i in range(0, n_features, self.feature_batch_size):\n            batch = self.feature_names[i:i+self.feature_batch_size]\n            feature_batches.append(batch)\n        \n        print(f\"Split {n_features} features into {len(feature_batches)} batches\")\n        print(f\"Total lag strategies (including odd/even): {len(self.lag_strategies)}\")\n        \n        # Calculate total lag values\n        total_lags = sum(len(lags) for lags in self.lag_strategies.values())\n        print(f\"Total unique lag values: {total_lags}\")\n        \n        # Train models for each combination of feature batch and lag strategy\n        total_models = len(feature_batches) * len(self.lag_strategies)\n        model_count = 0\n        \n        for strategy_name, lag_list in self.lag_strategies.items():\n            print(f\"\\nTraining {strategy_name} strategy (lags: {lag_list})\")\n            \n            for batch_idx, feature_batch in enumerate(feature_batches):\n                model_count += 1\n                print(f\"Progress: {model_count}/{total_models} models\")\n                \n                self.train_feature_batch(feature_batch, X, y, strategy_name, lag_list)\n                \n                # Clean up periodically\n                if batch_idx % 2 == 0:\n                    gc.collect()\n        \n        # Initialize equal weights\n        for model_name in self.models:\n            self.model_weights[model_name] = 1.0 / len(self.models)\n        \n        print(f\"\\nTotal models trained: {len(self.models)}\")\n        \n    def predict_batch(self, X, feature_batch, strategy_name, lag_list):\n        \"\"\"Make predictions for a specific feature batch and lag strategy.\"\"\"\n        model_name = f\"{strategy_name}_{feature_batch[0]}_{feature_batch[-1]}\"\n        \n        if model_name not in self.models:\n            return None\n            \n        # Select features\n        X_batch = X[feature_batch].copy()\n        \n        # Create lag features\n        X_with_lags = self.create_lag_features_batch(X_batch, lag_list)\n        \n        # Scale and predict\n        X_scaled = self.scalers[model_name].transform(X_with_lags)\n        predictions = self.models[model_name].predict(X_scaled)\n        \n        # Clean up\n        del X_batch, X_with_lags, X_scaled\n        gc.collect()\n        \n        return predictions\n    \n    def predict(self, X):\n        \"\"\"Make ensemble predictions.\"\"\"\n        all_predictions = []\n        weights = []\n        \n        # Recreate feature batches\n        n_features = len(self.feature_names)\n        feature_batches = []\n        for i in range(0, n_features, self.feature_batch_size):\n            batch = self.feature_names[i:i+self.feature_batch_size]\n            if all(col in X.columns for col in batch):\n                feature_batches.append(batch)\n        \n        # Get predictions from each model\n        for strategy_name, lag_list in self.lag_strategies.items():\n            for feature_batch in feature_batches:\n                pred = self.predict_batch(X, feature_batch, strategy_name, lag_list)\n                if pred is not None:\n                    all_predictions.append(pred)\n                    model_name = f\"{strategy_name}_{feature_batch[0]}_{feature_batch[-1]}\"\n                    weights.append(self.model_weights.get(model_name, 1.0))\n        \n        # Weighted average\n        if all_predictions:\n            weights = np.array(weights) / np.sum(weights)\n            return np.average(all_predictions, axis=0, weights=weights)\n        else:\n            return np.zeros(len(X))\n\n# Main execution starts here\nprint(\"=\"*70)\nprint(\"ENHANCED CRYPTO PREDICTION WITH ADDITIONAL FEATURES AND ODD/EVEN LAGS\")\nprint(\"=\"*70)\n\n# Set pandas options\npd.options.mode.chained_assignment = None\npd.options.display.max_columns = None\n\n# Load training data\nprint(\"\\nLoading training data...\")\ntrain_df = pd.read_parquet('/kaggle/input/drw-crypto-market-prediction/train.parquet')\nprint(f\"Training data shape: {train_df.shape}\")\n\n# Optimize memory\ntrain_df = optimize_memory(train_df, verbose=True)\n\n# Extract labels\ny_train = train_df['label'].to_numpy().astype(np.float32)\n\n# Select features\nfeature_cols = get_feature_columns(train_df)\nX_train = train_df[feature_cols].copy()\n\n# Clean up\ndel train_df\ngc.collect()\n\n# Further optimize X_train\nX_train = optimize_memory(X_train, verbose=True)\n\n# Create and train the ensemble with more epochs\nensemble = IncrementalLagEnsemble(\n    feature_batch_size=35,  # Slightly smaller batches due to more features\n    lag_batch_size=5,\n    n_epochs=9  # Increased to 9 as requested\n)\nensemble.fit(X_train, y_train)\n\n# Clean up training data\ndel X_train, y_train\ngc.collect()\n\n# Load test data\nprint(\"\\n\" + \"=\"*50)\nprint(\"LOADING TEST DATA\")\nprint(\"=\"*50)\n\ntest_df = pd.read_parquet('/kaggle/input/drw-crypto-market-prediction/test.parquet')\nprint(f\"Test data shape: {test_df.shape}\")\n\n# Optimize memory\ntest_df = optimize_memory(test_df, verbose=True)\n\n# Timestamp reconstruction\ntimestamp_recon_path = '/kaggle/input/the-order-of-the-test-rows-2/closest_rows.csv'\nuse_timestamp_reconstruction = os.path.exists(timestamp_recon_path)\n\nif use_timestamp_reconstruction:\n    print(\"\\nApplying timestamp reconstruction...\")\n    \n    t = pd.Series(pd.read_csv(timestamp_recon_path)['0'].to_numpy())\n    print(f\"Timestamps loaded: {len(t)}\")\n    \n    # Process timestamps\n    t -= 10080\n    t[t < 0] = 538149\n    \n    t = t.sort_values()\n    t[t <= len(t)] = np.arange(t[t <= len(t)].shape[0])\n    t = t.sort_index()\n    \n    t = pd.Series(np.arange(538150), index=t.to_numpy()).sort_index()\n    \n    # Sort test data\n    test_df = test_df.iloc[t.to_numpy()]\n    print(\"Test data sorted by reconstructed timestamps\")\nelse:\n    print(\"No timestamp reconstruction file found\")\n    t = pd.Series(np.arange(len(test_df)))\n\n# Select same features as training\nX_test = test_df[feature_cols].copy()\ndel test_df\ngc.collect()\n\n# Optimize test features\nX_test = optimize_memory(X_test, verbose=True)\n\n# Make predictions in chunks\nprint(\"\\n\" + \"=\"*50)\nprint(\"MAKING PREDICTIONS\")\nprint(\"=\"*50)\n\nchunk_size = 15000  # Smaller chunks due to more features and lags\nn_samples = len(X_test)\ny_pred = np.zeros(n_samples, dtype=np.float32)\n\nn_chunks = (n_samples + chunk_size - 1) // chunk_size\nprint(f\"Processing {n_chunks} chunks of size {chunk_size}\")\n\nfor i in range(0, n_samples, chunk_size):\n    end_idx = min(i + chunk_size, n_samples)\n    chunk_num = i // chunk_size + 1\n    \n    if chunk_num % 10 == 0:\n        print(f\"\\nChunk {chunk_num}/{n_chunks} (rows {i}-{end_idx})\")\n    \n    # Get chunk\n    X_chunk = X_test.iloc[i:end_idx]\n    \n    # Predict\n    y_pred[i:end_idx] = ensemble.predict(X_chunk).astype(np.float32)\n    \n    # Clean up periodically\n    if chunk_num % 5 == 0:\n        gc.collect()\n\n# Clean up test data\ndel X_test\ngc.collect()\n\n# Display statistics\nprint(\"\\n\" + \"=\"*50)\nprint(\"PREDICTION STATISTICS\")\nprint(\"=\"*50)\n\npred_series = pd.Series(y_pred)\nprint(pred_series.describe())\n\n# Create visualizations\nfig, axes = plt.subplots(2, 2, figsize=(12, 8))\n\n# Cumulative sum\naxes[0, 0].plot(np.cumsum(y_pred))\naxes[0, 0].set_title('Cumulative Predictions')\naxes[0, 0].set_xlabel('Index')\naxes[0, 0].set_ylabel('Cumulative Sum')\naxes[0, 0].grid(True, alpha=0.3)\n\n# Distribution\naxes[0, 1].hist(y_pred, bins=50, alpha=0.7, edgecolor='black')\naxes[0, 1].set_title('Prediction Distribution')\naxes[0, 1].set_xlabel('Value')\naxes[0, 1].set_ylabel('Count')\n\n# First 2000 predictions\naxes[1, 0].plot(y_pred[:2000], alpha=0.7)\naxes[1, 0].set_title('First 2000 Predictions')\naxes[1, 0].set_xlabel('Index')\naxes[1, 0].set_ylabel('Prediction')\n\n# Rolling mean and std\nwindow = 1000\nrolling_mean = pred_series.rolling(window).mean()\nrolling_std = pred_series.rolling(window).std()\n\naxes[1, 1].plot(rolling_mean, label='Mean')\naxes[1, 1].fill_between(\n    range(len(rolling_mean)),\n    rolling_mean - rolling_std,\n    rolling_mean + rolling_std,\n    alpha=0.3,\n    label='±1 Std'\n)\naxes[1, 1].set_title(f'Rolling Statistics (window={window})')\naxes[1, 1].set_xlabel('Index')\naxes[1, 1].set_ylabel('Value')\naxes[1, 1].legend()\n\nplt.tight_layout()\nplt.show()\n\n# Additional analysis plot\nfig, ax = plt.subplots(1, 1, figsize=(10, 6))\n\n# Prediction volatility over time\nwindow_sizes = [100, 500, 1000, 5000]\nfor window in window_sizes:\n    rolling_vol = pred_series.rolling(window).std()\n    ax.plot(rolling_vol, label=f'Window {window}', alpha=0.7)\n\nax.set_title('Prediction Volatility Over Time')\nax.set_xlabel('Index')\nax.set_ylabel('Rolling Standard Deviation')\nax.legend()\nax.grid(True, alpha=0.3)\nplt.show()\n\n# Prepare submission\nprint(\"\\n\" + \"=\"*50)\nprint(\"PREPARING SUBMISSION\")\nprint(\"=\"*50)\n\nsubmission = pd.read_csv('/kaggle/input/drw-crypto-market-prediction/sample_submission.csv')\n\nif use_timestamp_reconstruction:\n    submission = submission.iloc[t.to_numpy()]\n    submission['prediction'] = y_pred\n    submission = submission.sort_index()\nelse:\n    submission['prediction'] = y_pred\n\n# Save submission\nsubmission.to_csv('submission.csv', index=False)\nprint(\"Submission saved to 'submission.csv'\")\n\n# Display submission info\nprint(\"\\nSubmission preview:\")\nprint(submission.head())\nprint(f\"\\nSubmission shape: {submission.shape}\")\nprint(f\"Prediction range: [{submission['prediction'].min():.6f}, {submission['prediction'].max():.6f}]\")\nprint(f\"Mean: {submission['prediction'].mean():.6f}\")\nprint(f\"Std: {submission['prediction'].std():.6f}\")\n\n# Percentile information\npercentiles = [1, 5, 10, 25, 50, 75, 90, 95, 99]\nprint(\"\\nPrediction percentiles:\")\nfor p in percentiles:\n    value = np.percentile(submission['prediction'], p)\n    print(f\"  {p}th percentile: {value:.6f}\")\n\nprint(\"\\n\" + \"=\"*50)\nprint(\"PROCESS COMPLETED SUCCESSFULLY!\")\nprint(\"=\"*50)\n\n# Summary\nprint(\"\\nModel Summary:\")\nprint(f\"- Total features used: {len(feature_cols)}\")\nprint(f\"- Feature batches: {len(feature_cols) // ensemble.feature_batch_size + 1}\")\nprint(f\"- Lag strategies (including odd/even): {len(ensemble.lag_strategies)}\")\nprint(f\"- Total models: {len(ensemble.models)}\")\nprint(f\"- Model type: SGDRegressor with warm start\")\nprint(f\"- Training: Incremental with partial_fit ({ensemble.n_epochs} epochs)\")\nprint(f\"- Prediction: Weighted ensemble average\")\n\n# Count total unique lag values\nall_lags = set()\nfor lags in ensemble.lag_strategies.values():\n    all_lags.update(lags)\nprint(f\"- Total unique lag values: {len(all_lags)}\")\n\nprint(\"\\nDone!\")","metadata":{"trusted":true},"outputs":[],"execution_count":null},{"cell_type":"code","source":"# #!/usr/bin/env python\n# # coding: utf-8\n\n# # 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\n# import numpy as np # linear algebra\n# import pandas as pd # data processing, CSV file I/O (e.g. pd.read_csv)\n# import matplotlib.pyplot as plt\n# from sklearn.linear_model import Ridge, SGDRegressor\n# from sklearn.preprocessing import StandardScaler\n# from sklearn.ensemble import VotingRegressor\n# import os\n# import gc\n# import warnings\n# from collections import defaultdict\n# warnings.filterwarnings('ignore')\n\n# # Input data files are available in the read-only \"../input/\" directory\n# for dirname, _, filenames in os.walk('/kaggle/input'):\n#     for filename in filenames:\n#         print(os.path.join(dirname, filename))\n\n# def optimize_memory(df, verbose=True):\n#     \"\"\"\n#     Optimize memory usage by downcasting numeric types where possible.\n#     \"\"\"\n#     if verbose:\n#         start_mem = df.memory_usage().sum() / 1024**2\n#         print(f'Memory usage before optimization: {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.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 optimization: {end_mem:.2f} MB')\n#         print(f'Decreased by {100 * (start_mem - end_mem) / start_mem:.1f}%')\n    \n#     return df\n\n# def create_advanced_features(df):\n#     \"\"\"Create advanced engineered features for better model performance.\"\"\"\n#     print(\"Creating advanced engineered features...\")\n    \n#     # Ensure we have the necessary columns\n#     required_cols = ['buy_qty', 'sell_qty', 'bid_qty', 'ask_qty', 'volume']\n#     missing_cols = [col for col in required_cols if col not in df.columns]\n    \n#     if missing_cols:\n#         print(f\"Warning: Missing columns for feature engineering: {missing_cols}\")\n#         return df\n    \n#     # Volume-based features\n#     if 'volume' in df.columns:\n#         df['log_volume_squared'] = np.log1p(df['volume']) ** 2\n#         df['volume_sqrt'] = np.sqrt(df['volume'])\n    \n#     # Order flow features\n#     if all(col in df.columns for col in ['buy_qty', 'sell_qty']):\n#         df['order_flow_ratio'] = df['buy_qty'] / (df['sell_qty'] + 1)\n#         df['order_flow_diff'] = df['buy_qty'] - df['sell_qty']\n#         df['order_flow_imbalance_squared'] = ((df['buy_qty'] - df['sell_qty']) / (df['buy_qty'] + df['sell_qty'] + 1)) ** 2\n    \n#     # Spread and liquidity features\n#     if all(col in df.columns for col in ['bid_qty', 'ask_qty']):\n#         df['depth_imbalance_ratio'] = (df['bid_qty'] - df['ask_qty']) / (df['bid_qty'] + df['ask_qty'] + 1)\n#         df['total_depth'] = df['bid_qty'] + df['ask_qty']\n#         df['log_total_depth'] = np.log1p(df['total_depth'])\n    \n#     # Momentum features (if price-related columns exist)\n#     price_cols = [col for col in df.columns if 'price' in col.lower() or col in ['X363', 'X405', 'X321']]\n#     if price_cols:\n#         for col in price_cols[:3]:  # Limit to top 3 price columns\n#             df[f'{col}_momentum'] = df[col].pct_change(5).fillna(0)\n#             df[f'{col}_volatility'] = df[col].rolling(10).std().fillna(0)\n    \n#     return df\n\n# def get_feature_columns(df):\n#     \"\"\"Get all feature columns including the extended list and new variables.\"\"\"\n#     # Extended feature list provided\n#     extended_features = [\n#         'X727', 'X427', 'X288', 'X721', 'X312', 'X421', 'X471', 'X573', 'X780', 'X255',\n#         'X144', 'X299', 'X301', 'X563', 'X737', 'X702', 'ask_qty', 'X507', 'X306', 'X501',\n#         'X303', 'amihud_illiquidity', 'X586', 'X43', 'X517', 'X248', 'X137', 'X757', 'X196',\n#         'X777', 'X280', 'X266', 'X689', 'X294', 'X492', 'X555', 'X731', 'X262', 'X576',\n#         'X13', 'X518', 'X502', 'X558', 'pin_proxy', 'X6', 'X602', 'X695', 'X703', 'X413',\n#         'X660', 'X37', 'X15', 'X310', 'X512', 'X362', 'X631', 'X214', 'X562', 'X488',\n#         'X510', 'X256', 'X35', 'X128', 'X86', 'X170', 'X30', 'X265', 'X323', 'X559',\n#         'X348', 'X130', 'X529', 'X20', 'X4', 'X90', 'X192', 'X91', 'X582', 'X99',\n#         'X24', 'X317', 'X707', 'X653', 'X519', 'X557', 'X371', 'X415', 'X84', 'X83',\n#         'order_toxicity', 'X360', 'X111', 'X699', 'X187', 'X591', 'X637', 'X567', 'X577',\n#         'X313', 'X60', 'X671', 'X698', 'X701', 'X725', 'X292', 'X638', 'X741', 'X379',\n#         'X700', 'X614', 'X676', 'X516', 'X697', 'X611', 'X311', 'X615', 'X706', 'X466',\n#         'X571', 'X451', 'X17', 'X584', 'X436', 'X305', 'liquidity_consumption', 'X34', 'X282',\n#         'X681', 'X7', 'X208', 'X41', 'X536', 'X548', 'X296', 'X776', 'X87', 'X40',\n#         'X570', 'X539', 'X474', 'X753', 'X425', 'X217', 'X199', 'X18', 'X609', 'X21',\n#         'X277', 'X279', 'X326', 'X540', 'X688', 'X553', 'X452', 'X738', 'X183', 'X759',\n#         'bid_ask_ratio', 'X495', 'volume_participation', 'X715', 'X385', 'X291', 'X409', 'X112',\n#         'X693', 'X102', 'X318', 'X705', 'X556', 'X547'\n#     ]\n    \n#     # Core features that are usually important\n#     core_features = [\n#         'X363', 'X405', 'X321', 'X175', 'X179', 'X197', 'X22', 'X181', 'X28', 'X169',\n#         'X198', 'X173', 'X338', 'X344', 'X587', 'X450', 'X97', 'X52', 'X444', 'X598',\n#         'X297', 'X138', 'X572', 'X343', 'X438', 'X459', 'X758', 'X25', 'buy_qty', \n#         'sell_qty', 'volume', 'bid_qty'\n#     ]\n    \n#     # NEW features to add\n#     new_features = [\n#         'X363', 'X321', 'X405', 'X730', 'X523', 'X756', 'X589', 'X462', 'X779', 'log_liquidity',\n#         'X25', 'X532', 'X520', 'X329', 'X383', 'X751', 'X535', 'X639', 'X596', 'X761',\n#         'X145', 'X709', 'X173', 'X245', 'X168', 'X171', 'X241', 'X31', 'X105', 'X63',\n#         'X263', 'X426', 'X286', 'X357', 'X399', 'X315', 'X468', 'X131', 'X647', 'log_spread',\n#         'X752', 'X254', 'X592', 'X733', 'X636', 'X394', 'X527', 'X180', 'X367', 'X38',\n#         'X634', 'X718', 'X387', 'X429', 'X345', 'X344', 'X253', 'X469', 'X446', 'X125',\n#         'X760', 'X186', 'X711', 'X150', 'X661', 'X215', 'X403', 'X141', 'X771', 'X453',\n#         'X401', 'X629', 'X616', 'X281', 'X432', 'X283', 'X244', 'X440', 'X430', 'X382',\n#         'X175', 'X95', 'X444', 'X189', 'X55', 'X605', 'X663', 'X194', 'X439', 'X670',\n#         'X483', 'X163', 'X376', 'X71', 'X650', 'X203', 'X8', 'X624', 'X160', 'X100',\n#         'X14', 'X511', 'X59', 'X302', 'X81', 'X325', 'X514', 'X649', 'X447', 'X538',\n#         'X443', 'X39', 'X343', 'X12', 'X678', 'X775', 'X498', 'X249', 'X42', 'X384',\n#         'kyle_lambda', 'X349', 'X356', 'X2', 'X250', 'X397', 'X685', 'X568', 'X136', 'X496',\n#         'X53', 'X66', 'X374', 'X590', 'X668', 'X585', 'X677', 'X667', 'X530', 'X28',\n#         'X64', 'X407', 'X494', 'X770', 'X710', 'X526', 'X644', 'X167', 'X190', 'X723',\n#         'X33', 'X579', 'X206'\n#     ]\n    \n#     # ADDITIONAL features requested by user (first batch)\n#     additional_features = [\n#         'X525', 'X267', 'X166', 'X719', 'X489', 'X758', 'X652', 'X433', 'X778', 'X428',\n#         'X617', 'X259', 'X633', 'X565', 'X364', 'depth_ratio', 'X550', 'X687', 'X610', 'X599',\n#         'X717', 'X587', 'X143', 'X506', 'X546', 'X505', 'X159', 'X574', 'X278', 'X458',\n#         'X1', 'X749', 'X155', 'X651', 'X470', 'X580', 'X445', 'X373', 'X82', 'X607',\n#         'X298', 'X221', 'X388', 'X120', 'X391', 'X23', 'X679', 'X377', 'X767', 'X755',\n#         'X566', 'X424', 'X438', 'X198', 'X300', 'X268', 'X434', 'X290', 'X368', 'X464',\n#         'X119', 'X197', 'X597', 'X157', 'X485', 'X127', 'X101', 'X533', 'X235', 'X712',\n#         'X154', 'X239', 'X10', 'X420', 'X449', 'X740', 'X227', 'X36', 'X358', 'X551',\n#         'X528', 'X285', 'X335', 'X152', 'X110', 'X68', 'X713', 'X402', 'X370', 'X735',\n#         'X200', 'X331', 'X473', 'X162', 'X213', 'X322', 'X289', 'X477', 'X113', 'X560',\n#         'X672', 'X621', 'X682', 'X5', 'X72', 'X44', 'X419', 'buy_pressure', 'X242', 'volume',\n#         'X472', 'X332', 'X441', 'buy_sell_ratio', 'pressure_ratio', 'X508', 'X594', 'X191',\n#         'X261', 'X603', 'net_pressure', 'order_flow_imbalance', 'sell_pressure', 'X240', 'X673',\n#         'X608', 'X509', 'X165', 'X720', 'X314', 'X522', 'X531', 'X625', 'bid_depth_ratio',\n#         'X435', 'X293', 'X486', 'price_efficiency', 'X716', 'X627', 'X626', 'X169', 'X613',\n#         'X680', 'X544', 'X115', 'X307', 'X665', 'X465', 'X347', 'X728', 'X70', 'log_volume',\n#         'X340', 'X459', 'X56', 'X395', 'X354', 'X51', 'X732', 'X247', 'X324', 'X316',\n#         'X76', 'X341', 'X739', 'X601', 'X386', 'X683', 'X149', 'X193', 'X628', 'X309',\n#         'X351', 'X393'\n#     ]\n    \n#     # ADDITIONAL features requested by user (second batch)\n#     additional_features_2 = [\n#         'X158', 'X116', 'X74', 'X237', 'X238', 'X640', 'X251', 'X185', 'X654', 'X178',\n#         'X748', 'X177', 'X26', 'X658', 'X754', 'X58', 'X521', 'X272', 'X499', 'X355',\n#         'X669', 'X220', 'X223', 'X258', 'X457', 'X85', 'aggressive_ratio', 'X772', 'X537', 'X89',\n#         'X389', 'X153', 'X412', 'X234', 'X478', 'X491', 'X32', 'X400', 'X334', 'X724',\n#         'X210', 'normalized_spread', 'liquidity_imbalance', 'X635', 'X656', 'X236', 'X88', 'X745',\n#         'X107', 'X148', 'X619', 'X304', 'X632', 'X390', 'depth_imbalance', 'X561', 'X588',\n#         'X106', 'X448', 'X480', 'X65', 'X747', 'X233', 'X271', 'X138', 'X600', 'X297',\n#         'X328', 'X22', 'total_liquidity', 'X350', 'X497', 'X179', 'X664', 'X174', 'X16',\n#         'X67', 'X69', 'X132', 'X47', 'X365', 'X541', 'X3', 'X414', 'X408', 'X359',\n#         'X121', 'X184', 'X646', 'X124', 'X552', 'X692', 'X228', 'X442', 'X736', 'X648',\n#         'X11', 'X114', 'X437', 'X396', 'X54', 'X612', 'X97', 'X484', 'X222', 'X493',\n#         'X542', 'X515', 'X243', 'X773', 'X135', 'X96', 'X431', 'X696', 'X353', 'liquidity_adjusted_volume',\n#         'ask_depth_ratio', 'X406', 'X684', 'X454', 'X417', 'X375', 'X333', 'X503', 'X109',\n#         'X691', 'X569', 'X742', 'X29', 'X418', 'X704', 'X729', 'X73', 'X226', 'X49',\n#         'X94', 'X774', 'X257', 'X161', 'X475', 'X308', 'X260', 'X479', 'X207', 'X19',\n#         'X337', 'X75', 'X118', 'X750', 'X9', 'X575', 'X202', 'X129', 'X232', 'X93',\n#         'X229', 'X404', 'X195', 'X581', 'X549', 'X246', 'X675', 'X411', 'X327', 'X369',\n#         'X765', 'X330', 'X46', 'X392', 'X504', 'X482', 'X450', 'X593', 'X181', 'X572',\n#         'X274', 'X714', 'X61', 'X554', 'X657', 'X45', 'X641', 'X320', 'X319', 'X662',\n#         'X744', 'X481', 'X766', 'X361', 'X78', 'X205', 'X270', 'X117', 'X487', 'X461',\n#         'bid_momentum', 'X722', 'X618', 'X108', 'X764'\n#     ]\n    \n#     # FINAL batch of features requested by user\n#     final_features = [\n#         'X346', 'X216', 'bid_ask_spread', 'X139', 'X230', 'X623', 'X769', 'X630', 'X416', 'X147',\n#         'X218', 'X666', 'X212', 'X336', 'X275', 'X62', 'X104', 'X146', 'X52', 'X410',\n#         'X476', 'X578', 'X524', 'X273', 'X743', 'X598', 'X79', 'X225', 'X734', 'X463',\n#         'X211', 'X252', 'X284', 'X133', 'X57', 'X352', 'buy_qty', 'X686', 'X456', 'X188',\n#         'bid_qty', 'X48', 'X209', 'ask_momentum', 'X287', 'X201', 'X674', 'X172', 'X564', 'X77',\n#         'X545', 'X620', 'X643', 'X342', 'X27', 'X366', 'X768', 'X645', 'X543', 'X595',\n#         'X295', 'X763', 'X490', 'X726', 'X398', 'X264', 'execution_quality', 'X762', 'X455', 'X655',\n#         'X534', 'X224', 'X156', 'X142', 'X123', 'X690', 'net_order_flow', 'X422', 'X694', 'X467',\n#         'X378', 'X606', 'sell_qty', 'X423', 'X381', 'X339', 'X126', 'X708', 'X500', 'X746',\n#         'X140', 'X182', 'X98', 'X164', 'X80', 'X122', 'X642', 'X276', 'X103', 'X338',\n#         'X372', 'X219', 'X151', 'X231', 'X460', 'X583', 'X380', 'X50', 'X513', 'X659',\n#         'X92', 'X176', 'X134', 'X269', 'X622', 'X204', 'X604'\n#     ]\n    \n#     # Combine all features and remove duplicates\n#     all_features = list(dict.fromkeys(\n#         extended_features + core_features + new_features + \n#         additional_features + additional_features_2 + final_features\n#     ))\n    \n#     # Filter to only available columns\n#     available_features = [col for col in all_features if col in df.columns]\n    \n#     # Add engineered features if they exist\n#     engineered_features = [\n#         'log_volume_squared', 'volume_sqrt', 'order_flow_ratio', 'order_flow_diff',\n#         'order_flow_imbalance_squared', 'depth_imbalance_ratio', 'total_depth', 'log_total_depth'\n#     ]\n    \n#     for feat in engineered_features:\n#         if feat in df.columns and feat not in available_features:\n#             available_features.append(feat)\n    \n#     # Add momentum and volatility features\n#     momentum_features = [col for col in df.columns if '_momentum' in col or '_volatility' in col]\n#     for feat in momentum_features:\n#         if feat not in available_features:\n#             available_features.append(feat)\n    \n#     print(f\"Found {len(available_features)} features out of {len(all_features)} requested\")\n    \n#     return available_features\n\n# class IncrementalLagEnsemble:\n#     \"\"\"Incrementally trains models on different lag configurations using SGD.\"\"\"\n    \n#     def __init__(self, feature_batch_size=20, lag_batch_size=5, n_epochs=20):\n#         self.feature_batch_size = feature_batch_size  # Reduced to 20 for memory efficiency\n#         self.lag_batch_size = lag_batch_size\n#         self.n_epochs = n_epochs  # Increased to 20 for better convergence\n#         self.models = {}\n#         self.scalers = {}\n#         self.feature_names = None\n#         self.model_weights = {}\n#         self.performance_history = defaultdict(list)\n        \n#         # Define lag strategies with odd/even splits for many configurations\n#         # Reduced number of strategies to focus on most important ones for memory efficiency\n#         self.lag_strategies = {\n#             # Micro lags - very short term\n#             'micro': [1, 2, 3, 4, 5],\n#             'micro_odd': [1, 3, 5, 7, 9],\n#             'micro_even': [2, 4, 6, 8, 10],\n            \n#             # Ultra short with odd/even\n#             'ultra_short': [6, 8, 10, 12, 15],\n#             'ultra_short_odd': [7, 9, 11, 13, 15, 17],\n#             'ultra_short_even': [6, 8, 10, 12, 14, 16],\n            \n#             # Short with odd/even\n#             'short': [20, 25, 30, 40, 50],\n#             'short_odd': [21, 25, 31, 41, 51],\n#             'short_even': [20, 24, 30, 40, 50],\n            \n#             # Short medium with odd/even\n#             'short_medium': [60, 75, 90, 105, 120],\n#             'short_medium_odd': [61, 75, 91, 105, 121],\n            \n#             # Medium with odd/even\n#             'medium': [150, 180, 210, 240, 300],\n#             'medium_odd': [151, 181, 211, 241, 301],\n            \n#             # Medium long with odd/even\n#             'medium_long': [360, 420, 480, 540, 600],\n#             'medium_long_odd': [361, 421, 481, 541, 601],\n            \n#             # Long with odd/even\n#             'long': [720, 840, 960, 1080, 1200],\n#             'long_odd': [721, 841, 961, 1081, 1201],\n            \n#             # Very long\n#             'very_long': [1440, 1800, 2160, 2520, 2880],\n            \n#             # Ultra long\n#             'ultra_long': [3600, 4320, 5040, 5760, 7200],\n            \n#             # Special patterns\n#             'hourly': [60, 120, 180, 240, 300, 360],  # Hourly patterns\n#             'daily': [1440, 2880, 4320, 5760, 7200],  # Daily patterns\n#         }\n        \n#     def create_lag_features_batch(self, df, lag_list):\n#         \"\"\"Create lag features for a batch of lags with memory optimization.\"\"\"\n#         lag_features = []\n        \n#         # Process lags in smaller groups to save memory\n#         for i in range(0, len(lag_list), 2):\n#             lag_group = lag_list[i:i+2]\n#             for lag in lag_group:\n#                 lagged = df.shift(-lag)\n#                 lagged.columns = [f'{col}_lag_{lag}' for col in df.columns]\n#                 lag_features.append(lagged)\n            \n#             # Garbage collect after each group\n#             if i % 4 == 0:\n#                 gc.collect()\n        \n#         result = pd.concat([df] + lag_features, axis=1)\n#         result = result.fillna(0)\n        \n#         # Clean up\n#         del lag_features\n#         gc.collect()\n        \n#         return result\n    \n#     def train_sgd_model(self, X, y, model_name):\n#         \"\"\"Train SGD model incrementally with adaptive learning and regularization.\"\"\"\n#         print(f\"  Training SGD model: {model_name} ({self.n_epochs} epochs)\")\n        \n#         # Initialize model and scaler if not exists\n#         if model_name not in self.models:\n#             # Enhanced SGD configuration for better convergence\n#             self.models[model_name] = SGDRegressor(\n#                 loss='huber',\n#                 penalty='elasticnet',\n#                 alpha=0.00005,  # Reduced regularization\n#                 l1_ratio=0.15,\n#                 learning_rate='adaptive',  # Changed to adaptive\n#                 eta0=0.015,  # Slightly higher initial learning rate\n#                 power_t=0.3,  # Adjusted decay\n#                 random_state=42,\n#                 warm_start=True,\n#                 max_iter=1500,  # Increased iterations\n#                 tol=5e-4,  # Slightly looser tolerance\n#                 epsilon=0.1,  # Huber loss threshold\n#                 average=True  # Use averaged SGD\n#             )\n#             self.scalers[model_name] = StandardScaler()\n            \n#         model = self.models[model_name]\n#         scaler = self.scalers[model_name]\n        \n#         # Train in epochs with smaller chunks for better convergence\n#         chunk_size = 10000  # Reduced to 10000 for memory efficiency\n#         for epoch in range(self.n_epochs):\n#             # Shuffle indices for each epoch\n#             indices = np.random.permutation(len(X))\n            \n#             # Learning rate decay schedule\n#             if epoch > 0 and epoch % 5 == 0:\n#                 model.set_params(eta0=model.eta0 * 0.9)\n            \n#             for start_idx in range(0, len(X), chunk_size):\n#                 end_idx = min(start_idx + chunk_size, len(X))\n                \n#                 # Get shuffled chunk\n#                 chunk_indices = indices[start_idx:end_idx]\n#                 X_chunk = X.iloc[chunk_indices]\n#                 y_chunk = y[chunk_indices]\n                \n#                 # Scale\n#                 if start_idx == 0 and epoch == 0:\n#                     X_scaled = scaler.fit_transform(X_chunk)\n#                 else:\n#                     X_scaled = scaler.transform(X_chunk)\n                \n#                 # Partial fit\n#                 model.partial_fit(X_scaled, y_chunk)\n                \n#                 # Aggressive garbage collection\n#                 if start_idx % (chunk_size * 2) == 0:\n#                     gc.collect()\n            \n#             # Print progress\n#             if epoch % 4 == 0:\n#                 print(f\"    Epoch {epoch+1}/{self.n_epochs} completed\")\n        \n#         return model\n    \n#     def train_feature_batch(self, feature_batch, X_full, y, strategy_name, lag_list):\n#         \"\"\"Train on a batch of features with specific lags.\"\"\"\n#         print(f\"\\n  Processing feature batch ({len(feature_batch)} features) with {strategy_name} lags\")\n        \n#         # Select feature batch\n#         X_batch = X_full[feature_batch].copy()\n        \n#         # Create lag features\n#         X_with_lags = self.create_lag_features_batch(X_batch, lag_list)\n        \n#         # Train SGD model\n#         model_name = f\"{strategy_name}_{feature_batch[0]}_{feature_batch[-1]}\"\n#         self.train_sgd_model(X_with_lags, y, model_name)\n        \n#         # Clean up\n#         del X_batch, X_with_lags\n#         gc.collect()\n        \n#     def fit(self, X, y):\n#         \"\"\"Fit ensemble using incremental training.\"\"\"\n#         print(\"Training Enhanced Incremental Lag Ensemble...\")\n#         print(f\"Total epochs per model: {self.n_epochs}\")\n        \n#         self.feature_names = X.columns.tolist()\n#         n_features = len(self.feature_names)\n        \n#         # Split features into batches\n#         feature_batches = []\n#         for i in range(0, n_features, self.feature_batch_size):\n#             batch = self.feature_names[i:i+self.feature_batch_size]\n#             feature_batches.append(batch)\n        \n#         print(f\"Split {n_features} features into {len(feature_batches)} batches\")\n#         print(f\"Total lag strategies: {len(self.lag_strategies)}\")\n        \n#         # Calculate total lag values\n#         all_lags = set()\n#         for lags in self.lag_strategies.values():\n#             all_lags.update(lags)\n#         print(f\"Total unique lag values: {len(all_lags)}\")\n        \n#         # Train models for each combination of feature batch and lag strategy\n#         total_models = len(feature_batches) * len(self.lag_strategies)\n#         model_count = 0\n        \n#         for strategy_name, lag_list in self.lag_strategies.items():\n#             print(f\"\\nTraining {strategy_name} strategy (lags: {lag_list})\")\n            \n#             for batch_idx, feature_batch in enumerate(feature_batches):\n#                 model_count += 1\n#                 print(f\"Progress: {model_count}/{total_models} models\")\n                \n#                 self.train_feature_batch(feature_batch, X, y, strategy_name, lag_list)\n                \n#                 # Aggressive cleanup after every model\n#                 gc.collect()\n        \n#         # Initialize equal weights\n#         for model_name in self.models:\n#             self.model_weights[model_name] = 1.0 / len(self.models)\n        \n#         print(f\"\\nTotal models trained: {len(self.models)}\")\n        \n#     def predict_batch(self, X, feature_batch, strategy_name, lag_list):\n#         \"\"\"Make predictions for a specific feature batch and lag strategy.\"\"\"\n#         model_name = f\"{strategy_name}_{feature_batch[0]}_{feature_batch[-1]}\"\n        \n#         if model_name not in self.models:\n#             return None\n            \n#         # Select features\n#         X_batch = X[feature_batch].copy()\n        \n#         # Create lag features\n#         X_with_lags = self.create_lag_features_batch(X_batch, lag_list)\n        \n#         # Scale and predict\n#         X_scaled = self.scalers[model_name].transform(X_with_lags)\n#         predictions = self.models[model_name].predict(X_scaled)\n        \n#         # Clean up\n#         del X_batch, X_with_lags, X_scaled\n#         gc.collect()\n        \n#         return predictions\n    \n#     def predict(self, X):\n#         \"\"\"Make ensemble predictions with weighted averaging.\"\"\"\n#         all_predictions = []\n#         weights = []\n        \n#         # Recreate feature batches\n#         n_features = len(self.feature_names)\n#         feature_batches = []\n#         for i in range(0, n_features, self.feature_batch_size):\n#             batch = self.feature_names[i:i+self.feature_batch_size]\n#             if all(col in X.columns for col in batch):\n#                 feature_batches.append(batch)\n        \n#         # Get predictions from each model\n#         prediction_count = 0\n#         for strategy_name, lag_list in self.lag_strategies.items():\n#             for feature_batch in feature_batches:\n#                 pred = self.predict_batch(X, feature_batch, strategy_name, lag_list)\n#                 if pred is not None:\n#                     all_predictions.append(pred)\n#                     model_name = f\"{strategy_name}_{feature_batch[0]}_{feature_batch[-1]}\"\n#                     weights.append(self.model_weights.get(model_name, 1.0))\n#                     prediction_count += 1\n                    \n#                     # Garbage collect after every 5 predictions\n#                     if prediction_count % 5 == 0:\n#                         gc.collect()\n        \n#         # Weighted average\n#         if all_predictions:\n#             weights = np.array(weights) / np.sum(weights)\n#             return np.average(all_predictions, axis=0, weights=weights)\n#         else:\n#             return np.zeros(len(X))\n\n# # Main execution starts here\n# print(\"=\"*70)\n# print(\"ULTRA-ENHANCED CRYPTO PREDICTION WITH COMPLETE FEATURE SET\")\n# print(\"=\"*70)\n\n# # Set pandas options\n# pd.options.mode.chained_assignment = None\n# pd.options.display.max_columns = None\n\n# # Load training data\n# print(\"\\nLoading training data...\")\n# train_df = pd.read_parquet('/kaggle/input/drw-crypto-market-prediction/train.parquet')\n# print(f\"Training data shape: {train_df.shape}\")\n\n# # Optimize memory\n# train_df = optimize_memory(train_df, verbose=True)\n\n# # Create advanced features\n# train_df = create_advanced_features(train_df)\n\n# # Extract labels\n# y_train = train_df['label'].to_numpy().astype(np.float32)\n\n# # Select features\n# feature_cols = get_feature_columns(train_df)\n# X_train = train_df[feature_cols].copy()\n\n# # Clean up\n# del train_df\n# gc.collect()\n\n# # Further optimize X_train\n# X_train = optimize_memory(X_train, verbose=True)\n\n# # Create and train the ensemble with more epochs\n# ensemble = IncrementalLagEnsemble(\n#     feature_batch_size=20,  # Reduced to 20 for memory efficiency\n#     lag_batch_size=5,\n#     n_epochs=20  # Increased to 20 for better convergence\n# )\n# ensemble.fit(X_train, y_train)\n\n# # Clean up training data\n# del X_train, y_train\n# gc.collect()\n\n# # Load test data\n# print(\"\\n\" + \"=\"*50)\n# print(\"LOADING TEST DATA\")\n# print(\"=\"*50)\n\n# test_df = pd.read_parquet('/kaggle/input/drw-crypto-market-prediction/test.parquet')\n# print(f\"Test data shape: {test_df.shape}\")\n\n# # Optimize memory\n# test_df = optimize_memory(test_df, verbose=True)\n\n# # Create advanced features for test data\n# test_df = create_advanced_features(test_df)\n\n# # Timestamp reconstruction\n# timestamp_recon_path = '/kaggle/input/the-order-of-the-test-rows-2/closest_rows.csv'\n# use_timestamp_reconstruction = os.path.exists(timestamp_recon_path)\n\n# if use_timestamp_reconstruction:\n#     print(\"\\nApplying timestamp reconstruction...\")\n    \n#     t = pd.Series(pd.read_csv(timestamp_recon_path)['0'].to_numpy())\n#     print(f\"Timestamps loaded: {len(t)}\")\n    \n#     # Process timestamps\n#     t -= 10080\n#     t[t < 0] = 538149\n    \n#     t = t.sort_values()\n#     t[t <= len(t)] = np.arange(t[t <= len(t)].shape[0])\n#     t = t.sort_index()\n    \n#     t = pd.Series(np.arange(538150), index=t.to_numpy()).sort_index()\n    \n#     # Sort test data\n#     test_df = test_df.iloc[t.to_numpy()]\n#     print(\"Test data sorted by reconstructed timestamps\")\n# else:\n#     print(\"No timestamp reconstruction file found\")\n#     t = pd.Series(np.arange(len(test_df)))\n\n# # Select same features as training\n# X_test = test_df[feature_cols].copy()\n# del test_df\n# gc.collect()\n\n# # Optimize test features\n# X_test = optimize_memory(X_test, verbose=True)\n\n# # Make predictions in chunks\n# print(\"\\n\" + \"=\"*50)\n# print(\"MAKING PREDICTIONS\")\n# print(\"=\"*50)\n\n# chunk_size = 5000  # Reduced to 5000 for extreme memory efficiency\n# n_samples = len(X_test)\n# y_pred = np.zeros(n_samples, dtype=np.float32)\n\n# n_chunks = (n_samples + chunk_size - 1) // chunk_size\n# print(f\"Processing {n_chunks} chunks of size {chunk_size}\")\n\n# for i in range(0, n_samples, chunk_size):\n#     end_idx = min(i + chunk_size, n_samples)\n#     chunk_num = i // chunk_size + 1\n    \n#     if chunk_num % 20 == 0:\n#         print(f\"\\nChunk {chunk_num}/{n_chunks} (rows {i}-{end_idx})\")\n    \n#     # Get chunk\n#     X_chunk = X_test.iloc[i:end_idx]\n    \n#     # Predict\n#     y_pred[i:end_idx] = ensemble.predict(X_chunk).astype(np.float32)\n    \n#     # Aggressive cleanup after every chunk\n#     gc.collect()\n\n# # Clean up test data\n# del X_test\n# gc.collect()\n\n# # Display statistics\n# print(\"\\n\" + \"=\"*50)\n# print(\"PREDICTION STATISTICS\")\n# print(\"=\"*50)\n\n# pred_series = pd.Series(y_pred)\n# print(pred_series.describe())\n\n# # Create visualizations\n# fig, axes = plt.subplots(2, 2, figsize=(12, 8))\n\n# # Cumulative sum\n# axes[0, 0].plot(np.cumsum(y_pred))\n# axes[0, 0].set_title('Cumulative Predictions')\n# axes[0, 0].set_xlabel('Index')\n# axes[0, 0].set_ylabel('Cumulative Sum')\n# axes[0, 0].grid(True, alpha=0.3)\n\n# # Distribution\n# axes[0, 1].hist(y_pred, bins=50, alpha=0.7, edgecolor='black')\n# axes[0, 1].set_title('Prediction Distribution')\n# axes[0, 1].set_xlabel('Value')\n# axes[0, 1].set_ylabel('Count')\n\n# # First 2000 predictions\n# axes[1, 0].plot(y_pred[:2000], alpha=0.7)\n# axes[1, 0].set_title('First 2000 Predictions')\n# axes[1, 0].set_xlabel('Index')\n# axes[1, 0].set_ylabel('Prediction')\n\n# # Rolling mean and std\n# window = 1000\n# rolling_mean = pred_series.rolling(window).mean()\n# rolling_std = pred_series.rolling(window).std()\n\n# axes[1, 1].plot(rolling_mean, label='Mean')\n# axes[1, 1].fill_between(\n#     range(len(rolling_mean)),\n#     rolling_mean - rolling_std,\n#     rolling_mean + rolling_std,\n#     alpha=0.3,\n#     label='±1 Std'\n# )\n# axes[1, 1].set_title(f'Rolling Statistics (window={window})')\n# axes[1, 1].set_xlabel('Index')\n# axes[1, 1].set_ylabel('Value')\n# axes[1, 1].legend()\n\n# plt.tight_layout()\n# plt.show()\n\n# # Additional analysis plot\n# fig, ax = plt.subplots(1, 1, figsize=(10, 6))\n\n# # Prediction volatility over time\n# window_sizes = [100, 500, 1000, 5000]\n# for window in window_sizes:\n#     rolling_vol = pred_series.rolling(window).std()\n#     ax.plot(rolling_vol, label=f'Window {window}', alpha=0.7)\n\n# ax.set_title('Prediction Volatility Over Time')\n# ax.set_xlabel('Index')\n# ax.set_ylabel('Rolling Standard Deviation')\n# ax.legend()\n# ax.grid(True, alpha=0.3)\n# plt.show()\n\n# # Prepare submission\n# print(\"\\n\" + \"=\"*50)\n# print(\"PREPARING SUBMISSION\")\n# print(\"=\"*50)\n\n# submission = pd.read_csv('/kaggle/input/drw-crypto-market-prediction/sample_submission.csv')\n\n# if use_timestamp_reconstruction:\n#     submission = submission.iloc[t.to_numpy()]\n#     submission['prediction'] = y_pred\n#     submission = submission.sort_index()\n# else:\n#     submission['prediction'] = y_pred\n\n# # Save submission\n# submission.to_csv('submission.csv', index=False)\n# print(\"Submission saved to 'submission.csv'\")\n\n# # Display submission info\n# print(\"\\nSubmission preview:\")\n# print(submission.head())\n# print(f\"\\nSubmission shape: {submission.shape}\")\n# print(f\"Prediction range: [{submission['prediction'].min():.6f}, {submission['prediction'].max():.6f}]\")\n# print(f\"Mean: {submission['prediction'].mean():.6f}\")\n# print(f\"Std: {submission['prediction'].std():.6f}\")\n\n# # Percentile information\n# percentiles = [1, 5, 10, 25, 50, 75, 90, 95, 99]\n# print(\"\\nPrediction percentiles:\")\n# for p in percentiles:\n#     value = np.percentile(submission['prediction'], p)\n#     print(f\"  {p}th percentile: {value:.6f}\")\n\n# print(\"\\n\" + \"=\"*50)\n# print(\"PROCESS COMPLETED SUCCESSFULLY!\")\n# print(\"=\"*50)\n\n# # Summary\n# print(\"\\nModel Summary:\")\n# print(f\"- Total features used: {len(feature_cols)}\")\n# print(f\"- Feature batches: {len(feature_cols) // ensemble.feature_batch_size + 1}\")\n# print(f\"- Lag strategies: {len(ensemble.lag_strategies)}\")\n# print(f\"- Total models: {len(ensemble.models)}\")\n# print(f\"- Model type: Enhanced SGDRegressor with adaptive learning\")\n# print(f\"- Training: Incremental with partial_fit ({ensemble.n_epochs} epochs)\")\n# print(f\"- Prediction: Weighted ensemble average\")\n\n# # Count total unique lag values\n# all_lags = set()\n# for lags in ensemble.lag_strategies.values():\n#     all_lags.update(lags)\n# print(f\"- Total unique lag values: {len(all_lags)}\")\n\n# print(\"\\nEnhanced Features:\")\n# print(\"- Advanced feature engineering applied\")\n# print(\"- Complete feature set (700+ features)\")\n# print(\"- Improved SGD configuration with adaptive learning\")\n# print(\"- Enhanced lag strategies including special patterns\")\n\n# print(\"\\nDone!\")","metadata":{"trusted":true},"outputs":[],"execution_count":null}]}