{"metadata":{"kernelspec":{"language":"python","display_name":"Python 3","name":"python3"},"language_info":{"name":"python","version":"3.11.13","mimetype":"text/x-python","codemirror_mode":{"name":"ipython","version":3},"pygments_lexer":"ipython3","nbconvert_exporter":"python","file_extension":".py"},"kaggle":{"accelerator":"none","dataSources":[{"sourceId":96164,"databundleVersionId":12993472,"sourceType":"competition"},{"sourceId":12405352,"sourceType":"datasetVersion","datasetId":7823194},{"sourceId":249869065,"sourceType":"kernelVersion"}],"dockerImageVersionId":31089,"isInternetEnabled":true,"language":"python","sourceType":"notebook","isGpuEnabled":false}},"nbformat_minor":4,"nbformat":4,"cells":[{"cell_type":"code","source":"import 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    # EXPANDED features for more comprehensive coverage\n    expanded_features = [\n        'X3', 'X9', 'X11', 'X16', 'X19', 'X26', 'X27', 'X29', 'X32', 'X45', 'X46', 'X47',\n        'X48', 'X49', 'X50', 'X54', 'X57', 'X58', 'X61', 'X62', 'X65', 'X67', 'X69', 'X73',\n        'X74', 'X75', 'X77', 'X78', 'X79', 'X80', 'X85', 'X89', 'X92', 'X93', 'X94', 'X96',\n        'X98', 'X103', 'X104', 'X106', 'X107', 'X108', 'X109', 'X114', 'X116', 'X117', 'X118',\n        'X121', 'X122', 'X123', 'X124', 'X126', 'X129', 'X132', 'X133', 'X134', 'X135', 'X139',\n        'X140', 'X142', 'X146', 'X147', 'X148', 'X151', 'X153', 'X156', 'X158', 'X161', 'X164',\n        'X172', 'X174', 'X176', 'X177', 'X178', 'X182', 'X184', 'X185', 'X188', 'X195', 'X201',\n        'X202', 'X204', 'X205', 'X207', 'X209', 'X210', 'X211', 'X212', 'X216', 'X218', 'X219',\n        'X220', 'X222', 'X223', 'X224', 'X225', 'X226', 'X228', 'X229', 'X230', 'X231', 'X232',\n        'X233', 'X234', 'X236', 'X237', 'X238', 'X243', 'X246', 'X251', 'X252', 'X257', 'X258',\n        'X260', 'X264', 'X269', 'X270', 'X271', 'X272', 'X273', 'X275', 'X276', 'X284', 'X287',\n        'X295', 'X304', 'X308', 'X319', 'X320', 'X327', 'X328', 'X330', 'X333', 'X334', 'X336',\n        'X337', 'X339', 'X342', 'X346', 'X350', 'X352', 'X353', 'X355', 'X359', 'X361', 'X365',\n        'X366', 'X369', 'X372', 'X375', 'X378', 'X380', 'X381', 'X389', 'X390', 'X392', 'X396',\n        'X398', 'X400', 'X404', 'X406', 'X408', 'X410', 'X411', 'X412', 'X414', 'X416', 'X417',\n        'X418', 'X422', 'X423', 'X431', 'X437', 'X442', 'X448', 'X454', 'X455', 'X456', 'X457',\n        'X460', 'X461', 'X463', 'X467', 'X475', 'X476', 'X478', 'X479', 'X480', 'X481', 'X482',\n        'X484', 'X487', 'X490', 'X491', 'X493', 'X497', 'X499', 'X500', 'X503', 'X504', 'X513',\n        'X515', 'X521', 'X524', 'X534', 'X537', 'X541', 'X542', 'X543', 'X545', 'X549', 'X552',\n        'X554', 'X561', 'X564', 'X569', 'X575', 'X578', 'X581', 'X583', 'X588', 'X593', 'X595',\n        'X600', 'X604', 'X606', 'X612', 'X618', 'X619', 'X620', 'X622', 'X623', 'X630', 'X632',\n        'X635', 'X640', 'X641', 'X642', 'X643', 'X645', 'X646', 'X648', 'X654', 'X655', 'X656',\n        'X657', 'X658', 'X659', 'X662', 'X664', 'X666', 'X669', 'X674', 'X675', 'X684', 'X686',\n        'X690', 'X691', 'X692', 'X694', 'X696', 'X704', 'X708', 'X714', 'X722', 'X724', 'X726',\n        'X729', 'X734', 'X736', 'X742', 'X743', 'X744', 'X745', 'X746', 'X747', 'X748', 'X750',\n        'X754', 'X762', 'X763', 'X764', 'X765', 'X766', 'X768', 'X769', 'X772', 'X773', 'X774',\n        'X781', 'X782', 'X783', 'X784', 'X785'\n    ]\n    \n    # Combine all features and remove duplicates\n    all_features = list(dict.fromkeys(extended_features + core_features + new_features + \n                                     additional_features + expanded_features))\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=6):\n        self.feature_batch_size = feature_batch_size\n        self.lag_batch_size = lag_batch_size\n        self.n_epochs = n_epochs  # Reduced from 9 to 6\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 fewer 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 = 30000  # Slightly larger chunks for faster training\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 % 2 == 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(\"=\"*80)\nprint(\"ENHANCED CRYPTO PREDICTION WITH MORE FEATURES AND OPTIMIZED EPOCHS\")\nprint(\"=\"*80)\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 fewer epochs but more features\nensemble = IncrementalLagEnsemble(\n    feature_batch_size=50,  # Increased batch size to handle more features efficiently\n    lag_batch_size=5,\n    n_epochs=6  # Reduced from 9 to 6 epochs\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 = 20000  # Increased chunk size for faster processing\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,"execution":{"iopub.status.busy":"2025-07-21T00:12:16.712969Z","iopub.execute_input":"2025-07-21T00:12:16.713339Z"}},"outputs":[],"execution_count":null}]}