{"metadata":{"kernelspec":{"language":"python","display_name":"Python 3","name":"python3"},"language_info":{"name":"python","version":"3.10.12","mimetype":"text/x-python","codemirror_mode":{"name":"ipython","version":3},"pygments_lexer":"ipython3","nbconvert_exporter":"python","file_extension":".py"},"kaggle":{"accelerator":"none","dataSources":[{"sourceId":84493,"databundleVersionId":9871156,"sourceType":"competition"}],"dockerImageVersionId":30822,"isInternetEnabled":true,"language":"python","sourceType":"notebook","isGpuEnabled":false}},"nbformat_minor":4,"nbformat":4,"cells":[{"cell_type":"code","source":"# Install all required libraries at once\n!pip install pytorch-tabnet fastFM dask joblib tqdm --quiet\n\n# Import necessary libraries\nimport os\nimport glob\nimport pandas as pd\nimport numpy as np\nimport dask.dataframe as dd\nimport gc\nimport psutil\nfrom joblib import Parallel, delayed\nfrom tqdm import tqdm\nfrom collections import Counter\nimport matplotlib.pyplot as plt\nimport seaborn as sns\n","metadata":{"_uuid":"8f2839f25d086af736a60e9eeb907d3b93b6e0e5","_cell_guid":"b1076dfc-b9ad-4769-8c92-a6c4dae69d19","trusted":true,"execution":{"iopub.status.busy":"2024-12-30T07:26:25.804442Z","iopub.execute_input":"2024-12-30T07:26:25.804849Z","iopub.status.idle":"2024-12-30T07:26:51.137268Z","shell.execute_reply.started":"2024-12-30T07:26:25.804814Z","shell.execute_reply":"2024-12-30T07:26:51.136110Z"}},"outputs":[],"execution_count":null},{"cell_type":"code","source":"# Install required libraries\n!pip install pytorch-tabnet fastFM dask\n\n# Import libraries\nimport os\nimport pandas as pd\nimport numpy as np\nimport dask.dataframe as dd\nfrom dask import delayed, compute\nfrom tqdm import tqdm\nimport gc\nimport psutil\n\n# Function to monitor memory usage\ndef print_memory_usage():\n    process = psutil.Process(os.getpid())\n    mem = process.memory_info().rss / 1e6  # in MB\n    print(f\"Current memory usage: {mem:.2f} MB\")\n\n# Define the sampling function\ndef sample_group(group, fraction=0.33):\n    N = len(group)\n    S = int(np.floor(fraction * N))\n    s = int(np.floor(S / 3))\n    \n    if S < 3 or s < 1:\n        if N >= 2:\n            start = group.iloc[[0]]\n            end = group.iloc[[-1]]\n            sampled = pd.concat([start, end], ignore_index=True)\n            return sampled\n        elif N == 1:\n            return group.iloc[[0]]\n        else:\n            return pd.DataFrame()\n    else:\n        start = group.iloc[:s]\n        mid_start = (N // 2) - (s // 2)\n        mid_end = mid_start + s\n        mid = group.iloc[mid_start:mid_end]\n        end = group.iloc[-s:]\n        sampled = pd.concat([start, mid, end], ignore_index=True)\n        return sampled\n\n# Define the delayed function\n@delayed\ndef process_partition(partition_id, partition_dir, sampled_dir, fraction=0.33):\n    partition_path = os.path.join(partition_dir, partition_id, 'part-0.parquet')\n    \n    if not os.path.isfile(partition_path):\n        print(f\"\\nPartition ID {partition_id}: File '{partition_path}' does not exist.\")\n        return None\n    \n    print(f\"\\nProcessing file: {partition_path}\")\n    \n    try:\n        df = pd.read_parquet(partition_path)\n    except Exception as e:\n        print(f\"Error reading {partition_path}: {e}\")\n        return None\n    \n    print(f\"Original Partition Shape: {df.shape}\")\n    \n    grouped = df.groupby(['symbol_id', 'date_id'])\n    \n    sampled_partition = grouped.apply(lambda grp: sample_group(grp, fraction=fraction)).reset_index(drop=True)\n    \n    print(f\"Sampled Partition Shape: {sampled_partition.shape}\")\n    \n    sampled_file_name = f'sampled_{partition_id}.parquet'\n    sampled_file_path = os.path.join(sampled_dir, sampled_file_name)\n    \n    try:\n        sampled_partition.to_parquet(sampled_file_path, index=False, compression='snappy')\n        print(f\"Sampled partition saved to '{sampled_file_path}'.\")\n    except Exception as e:\n        print(f\"Error saving sampled partition {sampled_file_path}: {e}\")\n        return None\n    \n    del df\n    del sampled_partition\n    gc.collect()\n    \n    return sampled_file_path\n\n# Create directory for sampled partitions\nsampled_dir = '/kaggle/working/sampled_partitions_dask/'\nos.makedirs(sampled_dir, exist_ok=True)\n\n# List all partition directories\npartition_dir = '/kaggle/input/jane-street-real-time-market-data-forecasting/train.parquet/'\npartition_ids = [d for d in os.listdir(partition_dir) if d.startswith('partition_id=')]\n\nprint(f\"Found {len(partition_ids)} partition directories.\")\n\n# Create and execute delayed tasks\ntasks = [process_partition(partition_id, partition_dir, sampled_dir, fraction=0.33) for partition_id in partition_ids]\nprint(\"Starting partition processing with Dask...\")\nresults = compute(*tasks)\n\n# Collect and verify sampled files\nsampled_files = [res for res in results if res is not None]\nprint(f\"\\nTotal sampled partition files saved: {len(sampled_files)}\")\n\n# Concatenate all sampled partitions into a final Dask DataFrame\nif sampled_files:\n    try:\n        ddf_sampled = dd.read_parquet(sampled_files, engine='pyarrow')\n    except Exception as e:\n        print(f\"Error reading sampled parquet files: {e}\")\n        ddf_sampled = None\n\n    if ddf_sampled is not None:\n        try:\n            num_rows = ddf_sampled.shape[0].compute()\n        except Exception as e:\n            print(f\"Error computing number of rows: {e}\")\n            num_rows = 'Unknown'\n    \n        num_cols = len(ddf_sampled.columns)\n    \n        print(f\"Sampled Dask DataFrame Shape: ({num_rows}, {num_cols})\")\n    \n        # Save the final sampled data as a single Parquet file\n        try:\n            ddf_sampled.to_parquet('/kaggle/working/final_sampled_data_dask.parquet', write_index=False, compression='snappy')\n            print(\"Final sampled data saved to '/kaggle/working/final_sampled_data_dask.parquet'.\")\n        except Exception as e:\n            print(f\"Error saving final sampled data: {e}\")\n    else:\n        print(\"Dask DataFrame creation failed.\")\nelse:\n    print(\"No sampled data to concatenate.\")\n\n# Monitor memory usage\nprint_memory_usage()\n\n# Clean up\ndel tasks\ndel results\ngc.collect()\n","metadata":{"trusted":true,"execution":{"iopub.status.busy":"2024-12-30T07:26:51.138868Z","iopub.execute_input":"2024-12-30T07:26:51.139521Z","iopub.status.idle":"2024-12-30T07:30:58.079011Z","shell.execute_reply.started":"2024-12-30T07:26:51.139476Z","shell.execute_reply":"2024-12-30T07:30:58.076255Z"}},"outputs":[],"execution_count":null},{"cell_type":"code","source":"\nimport shutil\nimport os\n\n# Path to the directory to delete\ndirectory_path = '/kaggle/working/sampled_partitions_dask'\n#directory_path = '/kaggle/working/imputed_partitions'\n# Check if the directory exists\nif os.path.exists(directory_path):\n    try:\n        # Delete the directory and its contents\n        shutil.rmtree(directory_path)\n        print(f\"Directory '{directory_path}' and its contents have been deleted.\")\n    except Exception as e:\n        print(f\"Error deleting directory '{directory_path}': {e}\")\nelse:\n    print(f\"Directory '{directory_path}' does not exist.\")\n","metadata":{"trusted":true,"execution":{"iopub.status.busy":"2024-12-30T07:36:35.201890Z","iopub.execute_input":"2024-12-30T07:36:35.202463Z","iopub.status.idle":"2024-12-30T07:36:35.210402Z","shell.execute_reply.started":"2024-12-30T07:36:35.202419Z","shell.execute_reply":"2024-12-30T07:36:35.209266Z"}},"outputs":[],"execution_count":null},{"cell_type":"code","source":"# ----------------------------- Preprocessing Pipeline with Efficient Validation -----------------------------","metadata":{"trusted":true},"outputs":[],"execution_count":null},{"cell_type":"markdown","source":"Install and Import Required Libraries","metadata":{}},{"cell_type":"code","source":"# Install required libraries\n!pip install pytorch-tabnet fastFM dask joblib tqdm --quiet\n\n# Import necessary libraries\nimport os\nimport glob\nimport pandas as pd\nimport numpy as np\nimport dask.dataframe as dd\nimport gc\nimport psutil\nfrom joblib import Parallel, delayed\nfrom tqdm import tqdm\nfrom collections import Counter\nimport matplotlib.pyplot as plt\nimport seaborn as sns\n","metadata":{"trusted":true,"execution":{"iopub.status.busy":"2024-12-30T07:37:10.151319Z","iopub.execute_input":"2024-12-30T07:37:10.151688Z","iopub.status.idle":"2024-12-30T07:37:15.550938Z","shell.execute_reply.started":"2024-12-30T07:37:10.151657Z","shell.execute_reply":"2024-12-30T07:37:15.549110Z"}},"outputs":[],"execution_count":null},{"cell_type":"markdown","source":"b. Define Feature Lists","metadata":{}},{"cell_type":"code","source":"# Define Top 30 Features\ntop_30_features = [\n    'feature_61', 'feature_20', 'feature_24', 'feature_25', 'feature_23',\n    'feature_29', 'feature_22', 'feature_30', 'feature_38', 'feature_04',\n    'feature_05', 'feature_26', 'feature_28', 'feature_07', 'feature_01',\n    'feature_08', 'feature_21', 'feature_27', 'feature_31', 'feature_47',\n    'feature_58', 'feature_62', 'feature_37', 'feature_36', 'feature_60',\n    'feature_17', 'feature_64', 'feature_33', 'feature_49', 'feature_15'\n]\n\n# Verify the list\nprint(\"Top 30 Feature Names:\")\nprint(top_30_features)\n\n# Define Features to Impute with Median and Mode\nmedian_features = ['feature_26', 'feature_21', 'feature_27', 'feature_31']\nmode_features = ['feature_04', 'feature_01']\n\nprint(\"\\nMedian Features:\")\nprint(median_features)\nprint(\"\\nMode Features:\")\nprint(mode_features)\n","metadata":{"trusted":true,"execution":{"iopub.status.busy":"2024-12-30T07:37:21.550569Z","iopub.execute_input":"2024-12-30T07:37:21.550990Z","iopub.status.idle":"2024-12-30T07:37:21.559651Z","shell.execute_reply.started":"2024-12-30T07:37:21.550926Z","shell.execute_reply":"2024-12-30T07:37:21.558468Z"}},"outputs":[],"execution_count":null},{"cell_type":"markdown","source":"c. Helper Functions","metadata":{}},{"cell_type":"code","source":"# Function to monitor memory usage\ndef print_memory_usage():\n    \"\"\"Print the current memory usage of the process.\"\"\"\n    process = psutil.Process(os.getpid())\n    mem = process.memory_info().rss / 1e6  # Convert bytes to MB\n    print(f\"Current memory usage: {mem:.2f} MB\")\n\n# Function to impute missing values\ndef impute_missing_values_vectorized(df, feature, window=7):\n    \"\"\"\n    Impute missing values in a feature using the rolling window median,\n    then fallback to forward/backward fill and global median.\n\n    Parameters:\n    - df (pd.DataFrame): DataFrame sorted by 'symbol_id' and 'time_id'.\n    - feature (str): The feature column to impute.\n    - window (int): The size of the rolling window.\n\n    Returns:\n    - pd.DataFrame: DataFrame with imputed values for the specified feature.\n    \"\"\"\n    # Compute rolling window median\n    rolling_median = df.groupby('symbol_id')[feature].transform(\n        lambda x: x.rolling(window=window, center=True, min_periods=1).median()\n    )\n\n    # Fill nulls with rolling window median\n    df[feature] = df[feature].fillna(rolling_median)\n\n    # Identify remaining nulls\n    remaining_nulls = df[feature].isnull()\n\n    if remaining_nulls.any():\n        # Forward fill within each group\n        df[feature] = df.groupby('symbol_id')[feature].transform(lambda x: x.ffill())\n\n        # Backward fill within each group\n        df[feature] = df.groupby('symbol_id')[feature].transform(lambda x: x.bfill())\n\n        # Fill any remaining nulls with global median\n        global_median = df[feature].median()\n        df[feature] = df[feature].fillna(global_median)\n\n    return df\n","metadata":{"trusted":true,"execution":{"iopub.status.busy":"2024-12-30T07:37:30.208542Z","iopub.execute_input":"2024-12-30T07:37:30.208910Z","iopub.status.idle":"2024-12-30T07:37:30.217438Z","shell.execute_reply.started":"2024-12-30T07:37:30.208882Z","shell.execute_reply":"2024-12-30T07:37:30.215926Z"}},"outputs":[],"execution_count":null},{"cell_type":"markdown","source":"d. Phase 1: Processing Source Partitions\n\n","metadata":{}},{"cell_type":"code","source":"# ----------------------------- Phase 1: Processing Source Partitions -----------------------------\n\n# Update Source Partitions to Include 5-9\nsource_partitions = [\n    'part.3.parquet',\n    'part.4.parquet',\n    'part.5.parquet',\n    'part.6.parquet',\n    'part.7.parquet',\n    'part.8.parquet',\n    'part.9.parquet'\n]\n\n# Define the directory containing the sampled partitions\nsampled_dir = '/kaggle/working/final_sampled_data_dask.parquet/'  # Adjust if necessary\n\n# Define the output directory for imputed source partitions\nimputed_dir = '/kaggle/working/imputed_partitions/'\nos.makedirs(imputed_dir, exist_ok=True)\n\n# Function to Process a Single Partition\ndef process_partition_file(partition_file, features, output_dir, window=7):\n    \"\"\"\n    Process a single Parquet partition: load, verify, sort, impute, and save.\n\n    Parameters:\n    - partition_file (str): Path to the partition file.\n    - features (list): List of feature names to process.\n    - output_dir (str): Directory to save imputed partitions.\n    - window (int): Rolling window size for imputation.\n\n    Returns:\n    - str: Status message indicating success or failure.\n    \"\"\"\n    try:\n        print(f\"\\nProcessing partition: {partition_file}\")\n\n        # Load the partition into a Pandas DataFrame using PyArrow\n        df_partition = pd.read_parquet(partition_file, engine='pyarrow')\n        print(f\"Loaded partition with shape {df_partition.shape}\")\n\n        # Verify required columns\n        required_columns = ['symbol_id', 'time_id', 'date_id'] + features\n        missing_columns = set(required_columns) - set(df_partition.columns)\n        if missing_columns:\n            error_msg = f\"Missing columns in {partition_file}: {missing_columns}\"\n            print(error_msg)\n            return error_msg  # Return error message\n\n        # Sort the DataFrame\n        df_partition = df_partition.sort_values(['symbol_id', 'time_id']).reset_index(drop=True)\n        print(\"DataFrame sorted by 'symbol_id' and 'time_id'.\")\n\n        # Identify columns with nulls\n        null_features = [feature for feature in features if df_partition[feature].isnull().any()]\n        print(f\"Features with nulls: {null_features}\")\n\n        # Impute each null-containing feature with a progress bar\n        for feature in tqdm(null_features, desc=f\"Imputing {os.path.basename(partition_file)}\"):\n            df_partition = impute_missing_values_vectorized(df_partition, feature, window=window)\n\n        # Verify Imputation\n        remaining_nulls = df_partition[features].isnull().sum()\n        print(\"\\nRemaining nulls after imputation:\")\n        print(remaining_nulls)\n\n        # Save the imputed partition\n        imputed_partition_filename = os.path.basename(partition_file).replace('.parquet', '_imputed.parquet')\n        imputed_partition_path = os.path.join(output_dir, imputed_partition_filename)\n        os.makedirs(output_dir, exist_ok=True)  # Create directory if it doesn't exist\n\n        df_partition.to_parquet(imputed_partition_path, compression='snappy', index=False)\n        print(f\"Imputed partition saved to {imputed_partition_path}.\")\n\n        # Clean up memory\n        del df_partition\n        gc.collect()\n\n        # Monitor memory usage\n        print_memory_usage()\n\n        return f\"Successfully processed {partition_file}\"\n\n    except Exception as e:\n        error_msg = f\"Error processing {partition_file}: {e}\"\n        print(error_msg)\n        return error_msg\n\n# Function to Process and Impute Partitions in Parallel\ndef process_and_impute_partitions(partition_files, features, output_dir, window):\n    \"\"\"\n    Process and impute multiple partitions in parallel.\n\n    Parameters:\n    - partition_files (list): List of partition file paths.\n    - features (list): List of feature names to process.\n    - output_dir (str): Directory to save imputed partitions.\n    - window (int): Rolling window size for imputation.\n\n    Returns:\n    - list: List of status messages.\n    \"\"\"\n    try:\n        results = Parallel(n_jobs=3, verbose=10)(\n            delayed(process_partition_file)(\n                pf,\n                features,\n                output_dir,\n                window\n            ) for pf in partition_files\n        )\n        return results\n    except Exception as e:\n        print(f\"Error in parallel processing: {e}\")\n        return []\n\n# Execute Phase 1: Process Source Partitions\nprint(\"\\n--- Phase 1: Processing Source Partitions ---\")\n\n# Define the source partition paths\nsource_partition_paths = [os.path.join(sampled_dir, sp) for sp in source_partitions]\n\n# Process and impute source partitions\nsource_results = process_and_impute_partitions(\n    partition_files=source_partition_paths,\n    features=top_30_features,\n    output_dir=imputed_dir,\n    window=7\n)\n\n# Print summary of Phase 1 results\nprint(\"\\nPhase 1 Processing Results:\")\nfor res in source_results:\n    print(res)\n\nprint(\"\\n--- Phase 1 Completed ---\\n\")\n\n# Monitor memory usage\nprint_memory_usage()\n\n# Clean up\ndel source_results\ngc.collect()\n\n# ------------------------------ Phase 1 End ------------------------------\n","metadata":{"trusted":true,"execution":{"iopub.status.busy":"2024-12-30T07:38:35.346357Z","iopub.execute_input":"2024-12-30T07:38:35.346781Z","iopub.status.idle":"2024-12-30T07:40:05.495457Z","shell.execute_reply.started":"2024-12-30T07:38:35.346750Z","shell.execute_reply":"2024-12-30T07:40:05.494130Z"}},"outputs":[],"execution_count":null},{"cell_type":"markdown","source":"e. Phase 2: Calculating Symbol Statistics","metadata":{}},{"cell_type":"code","source":"# ----------------------------- Phase 2: Calculate Symbol Statistics -----------------------------\n\n# Function to Calculate Symbol Statistics\ndef calculate_symbol_statistics(source_files, median_features, mode_features, imputed_dir):\n    \"\"\"\n    Calculate per-symbol medians and modes from source partitions.\n\n    Parameters:\n    - source_files (list): List of source partition file names.\n    - median_features (list): Features to calculate medians.\n    - mode_features (list): Features to calculate modes.\n    - imputed_dir (str): Directory containing the source partitions.\n\n    Returns:\n    - symbol_medians (pd.DataFrame): DataFrame with per-symbol medians.\n    - symbol_modes (pd.DataFrame): DataFrame with per-symbol modes.\n    \"\"\"\n    # Calculate medians\n    symbol_medians = []\n    for file in source_files:\n        imputed_file = file.replace('.parquet', '_imputed.parquet')\n        file_path = os.path.join(imputed_dir, imputed_file)\n        if not os.path.exists(file_path):\n            print(f\"Imputed source file {file_path} does not exist. Please process it first.\")\n            continue\n        ddf = dd.read_parquet(file_path, columns=['symbol_id'] + median_features, engine='pyarrow')\n        df = ddf.compute()\n        medians = df.groupby('symbol_id').median().reset_index()\n        symbol_medians.append(medians)\n        del df, ddf\n        gc.collect()\n\n    if symbol_medians:\n        combined_medians = pd.concat(symbol_medians, ignore_index=True)\n        symbol_medians = combined_medians.groupby('symbol_id').median().reset_index()\n    else:\n        symbol_medians = pd.DataFrame()\n        print(\"No medians calculated.\")\n\n    # Calculate modes\n    symbol_modes = []\n    for file in source_files:\n        imputed_file = file.replace('.parquet', '_imputed.parquet')\n        file_path = os.path.join(imputed_dir, imputed_file)\n        if not os.path.exists(file_path):\n            print(f\"Imputed source file {file_path} does not exist. Please process it first.\")\n            continue\n        ddf = dd.read_parquet(file_path, columns=['symbol_id'] + mode_features, engine='pyarrow')\n        df = ddf.compute()\n        # Define a function to compute mode safely\n        def safe_mode(x):\n            mode_series = x.mode()\n            if not mode_series.empty:\n                return mode_series.iloc[0]\n            else:\n                return pd.NA\n        modes = df.groupby('symbol_id').agg(\n            {feature: safe_mode for feature in mode_features}\n        ).reset_index()\n        symbol_modes.append(modes)\n        del df, ddf\n        gc.collect()\n\n    if symbol_modes:\n        combined_modes = pd.concat(symbol_modes, ignore_index=True)\n        symbol_modes = combined_modes.groupby('symbol_id').agg(\n            {feature: safe_mode for feature in mode_features}\n        ).reset_index()\n    else:\n        symbol_modes = pd.DataFrame()\n        print(\"No modes calculated.\")\n\n    return symbol_medians, symbol_modes\n\n# Execute Phase 2: Calculate Symbol Statistics\nprint(\"\\n--- Phase 2: Calculating Symbol Statistics ---\")\n\n# Calculate Symbol Statistics\nsymbol_medians, symbol_modes = calculate_symbol_statistics(\n    source_files=source_partitions,\n    median_features=median_features,\n    mode_features=mode_features,\n    imputed_dir=imputed_dir\n)\n\n# Display Symbol Medians\nprint(\"\\nSymbol Medians:\")\nprint(symbol_medians.head())\nprint(symbol_medians.isnull().sum())\n\n# Display Symbol Modes\nprint(\"\\nSymbol Modes:\")\nprint(symbol_modes.head())\nprint(symbol_modes.isnull().sum())\n\nprint(\"\\n--- Phase 2 Completed ---\\n\")\n\n# Monitor memory usage\nprint_memory_usage()\n\n# Clean up\ndel symbol_medians, symbol_modes\ngc.collect()\n\n# ------------------------------ Phase 2 End ------------------------------\n","metadata":{"trusted":true,"execution":{"iopub.status.busy":"2024-12-30T07:43:36.077012Z","iopub.execute_input":"2024-12-30T07:43:36.077480Z","iopub.status.idle":"2024-12-30T07:43:46.030452Z","shell.execute_reply.started":"2024-12-30T07:43:36.077450Z","shell.execute_reply":"2024-12-30T07:43:46.029250Z"}},"outputs":[],"execution_count":null},{"cell_type":"markdown","source":"f. Phase 3: Processing Problematic Partitions","metadata":{}},{"cell_type":"code","source":"# ----------------------------- Phase 3: Process Problematic Partitions -----------------------------\n\n# Define Problematic Partitions and Their Respective Features\nproblematic_partitions = {\n    'part.0.parquet': ['feature_26', 'feature_21', 'feature_27', 'feature_31', 'feature_04', 'feature_01'],\n    'part.1.parquet': ['feature_26', 'feature_21', 'feature_27', 'feature_31'],\n    'part.2.parquet': ['feature_26', 'feature_21', 'feature_27', 'feature_31']\n}\n\n# Function to Calculate Global Fallbacks\ndef calculate_global_fallbacks(symbol_medians, symbol_modes, median_features, mode_features):\n    \"\"\"\n    Calculate global medians and modes for fallback imputation.\n\n    Parameters:\n    - symbol_medians (pd.DataFrame): DataFrame with per-symbol medians.\n    - symbol_modes (pd.DataFrame): DataFrame with per-symbol modes.\n    - median_features (list): Features for medians.\n    - mode_features (list): Features for modes.\n\n    Returns:\n    - global_medians (dict): Dictionary of global medians.\n    - global_modes (dict): Dictionary of global modes.\n    \"\"\"\n    global_medians = {}\n    for feature in median_features:\n        if feature in symbol_medians.columns and not symbol_medians[feature].isnull().all():\n            global_medians[feature] = symbol_medians[feature].median()\n        else:\n            global_medians[feature] = 0  # Replace with an appropriate default if needed\n            print(f\"Global median for {feature} set to 0.\")\n\n    global_modes = {}\n    for feature in mode_features:\n        if feature in symbol_modes.columns:\n            mode_series = symbol_modes[feature].mode()\n            if not mode_series.empty:\n                global_modes[feature] = mode_series.iloc[0]\n            else:\n                global_modes[feature] = 'Unknown'  # Replace with an appropriate default if needed\n                print(f\"Global mode for {feature} set to 'Unknown'.\")\n        else:\n            global_modes[feature] = 'Unknown'  # Replace with an appropriate default if needed\n            print(f\"Global mode for {feature} set to 'Unknown'.\")\n\n    return global_medians, global_modes\n\n# Function to Process a Single Partition with Enhanced Imputation\ndef process_partition_file_enhanced(partition_file, features, output_dir, window,\n                                    symbol_medians, symbol_modes, global_medians, global_modes):\n    \"\"\"\n    Process a single Parquet partition: load, verify, sort, impute, handle problematic features, and save.\n\n    Parameters:\n    - partition_file (str): Path to the partition file.\n    - features (list): List of feature names to process.\n    - output_dir (str): Directory to save imputed partitions.\n    - window (int): Rolling window size for imputation.\n    - symbol_medians (pd.DataFrame): DataFrame with per-symbol medians.\n    - symbol_modes (pd.DataFrame): DataFrame with per-symbol modes.\n    - global_medians (dict): Dictionary of global medians.\n    - global_modes (dict): Dictionary of global modes.\n\n    Returns:\n    - str: Status message indicating success or failure.\n    \"\"\"\n    try:\n        print(f\"\\nProcessing partition: {partition_file}\")\n\n        # Load the partition into a Pandas DataFrame using PyArrow\n        df_partition = pd.read_parquet(partition_file, engine='pyarrow')\n        print(f\"Loaded partition with shape {df_partition.shape}\")\n\n        # Verify required columns\n        required_columns = ['symbol_id', 'time_id', 'date_id'] + features\n        missing_columns = set(required_columns) - set(df_partition.columns)\n        if missing_columns:\n            error_msg = f\"Missing columns in {partition_file}: {missing_columns}\"\n            print(error_msg)\n            return error_msg  # Return error message\n\n        # Sort the DataFrame\n        df_partition = df_partition.sort_values(['symbol_id', 'time_id']).reset_index(drop=True)\n        print(\"DataFrame sorted by 'symbol_id' and 'time_id'.\")\n\n        # Identify columns with nulls\n        null_features = [feature for feature in features if df_partition[feature].isnull().any()]\n        print(f\"Features with nulls: {null_features}\")\n\n        # Impute each null-containing feature with a progress bar\n        for feature in tqdm(null_features, desc=f\"Imputing {os.path.basename(partition_file)}\"):\n            df_partition = impute_missing_values_vectorized(df_partition, feature, window=window)\n\n        print(\"Standard imputation with rolling window median completed.\")\n\n        # Check if this partition is problematic and has specific features to handle\n        partition_name = os.path.basename(partition_file)\n        if partition_name in problematic_partitions:\n            problematic_features = problematic_partitions[partition_name]\n            for feature in problematic_features:\n                # Check if all values are still null after standard imputation\n                if df_partition[feature].isnull().all():\n                    print(f\"All values are null for {feature} in {partition_file}. Applying additional imputation.\")\n\n                    if feature in median_features:\n                        # Check if median data is available\n                        if feature in symbol_medians.columns:\n                            # Rename the median feature to avoid column name conflicts\n                            df_partition = df_partition.merge(\n                                symbol_medians[['symbol_id', feature]].rename(columns={feature: f\"{feature}_median\"}),\n                                on='symbol_id',\n                                how='left'\n                            )\n                            median_col = f\"{feature}_median\"\n                            # Fill nulls with per-symbol median\n                            df_partition[feature].fillna(df_partition[median_col], inplace=True)\n                            # Fill remaining nulls with global median\n                            df_partition[feature].fillna(global_medians.get(feature, 0), inplace=True)\n                            print(f\"Imputed {feature} using per-symbol median and global median.\")\n                            # Drop median column\n                            df_partition.drop(columns=[median_col], inplace=True)\n                        else:\n                            print(f\"Median data for {feature} not available. Filling with global median.\")\n                            df_partition[feature].fillna(global_medians.get(feature, 0), inplace=True)\n\n                    elif feature in mode_features:\n                        # Check if mode data is available\n                        if feature in symbol_modes.columns:\n                            # Rename the mode feature to avoid column name conflicts\n                            df_partition = df_partition.merge(\n                                symbol_modes[['symbol_id', feature]].rename(columns={feature: f\"{feature}_mode\"}),\n                                on='symbol_id',\n                                how='left'\n                            )\n                            mode_col = f\"{feature}_mode\"\n                            # Fill nulls with per-symbol mode\n                            df_partition[feature].fillna(df_partition[mode_col], inplace=True)\n                            # Fill remaining nulls with global mode\n                            df_partition[feature].fillna(global_modes.get(feature, 'Unknown'), inplace=True)\n                            print(f\"Imputed {feature} using per-symbol mode and global mode.\")\n                            # Drop mode column\n                            df_partition.drop(columns=[mode_col], inplace=True)\n                        else:\n                            print(f\"Mode data for {feature} not available. Filling with global mode.\")\n                            df_partition[feature].fillna(global_modes.get(feature, 'Unknown'), inplace=True)\n\n        # Verify Imputation\n        remaining_nulls = df_partition[features].isnull().sum()\n        print(\"\\nRemaining nulls after imputation:\")\n        print(remaining_nulls)\n\n        # Save the imputed partition\n        imputed_partition_filename = os.path.basename(partition_file).replace('.parquet', '_imputed.parquet')\n        imputed_partition_path = os.path.join(output_dir, imputed_partition_filename)\n        os.makedirs(output_dir, exist_ok=True)  # Create directory if it doesn't exist\n\n        df_partition.to_parquet(imputed_partition_path, compression='snappy', index=False)\n        print(f\"Imputed partition saved to {imputed_partition_path}.\")\n\n        # Clean up memory\n        del df_partition\n        gc.collect()\n\n        # Monitor memory usage\n        print_memory_usage()\n\n        return f\"Successfully processed {partition_file}\"\n\n    except Exception as e:\n        error_msg = f\"Error processing {partition_file}: {e}\"\n        print(error_msg)\n        return error_msg\n\n# Function to Process and Impute Problematic Partitions in Parallel\ndef process_and_impute_problematic_partitions(partition_files, features, output_dir, window, symbol_medians, symbol_modes, global_medians, global_modes):\n    \"\"\"\n    Process and impute multiple problematic partitions in parallel.\n\n    Parameters:\n    - partition_files (list): List of partition file paths.\n    - features (list): List of feature names to process.\n    - output_dir (str): Directory to save imputed partitions.\n    - window (int): Rolling window size for imputation.\n    - symbol_medians (pd.DataFrame): DataFrame with per-symbol medians.\n    - symbol_modes (pd.DataFrame): DataFrame with per-symbol modes.\n    - global_medians (dict): Dictionary of global medians.\n    - global_modes (dict): Dictionary of global modes.\n\n    Returns:\n    - list: List of status messages.\n    \"\"\"\n    try:\n        results = Parallel(n_jobs=3, verbose=10)(\n            delayed(process_partition_file_enhanced)(\n                pf,\n                features,\n                output_dir,\n                window,\n                symbol_medians,\n                symbol_modes,\n                global_medians,\n                global_modes\n            ) for pf in partition_files\n        )\n        return results\n    except Exception as e:\n        print(f\"Error in parallel processing: {e}\")\n        return []\n\n# Execute Phase 3: Process Problematic Partitions\n\nprint(\"\\n--- Phase 3: Processing Problematic Partitions ---\")\n\n# Define the directory containing the sampled partitions\nsampled_dir = '/kaggle/working/final_sampled_data_dask.parquet/'  # Adjust if necessary\n\n# Define the output directory for imputed partitions (same as before)\noutput_dir = '/kaggle/working/imputed_partitions/'\n\n# List of problematic partition file paths\nproblematic_partition_files = [os.path.join(sampled_dir, p) for p in problematic_partitions.keys()]\n\n# Calculate Symbol Statistics\nsymbol_medians, symbol_modes = calculate_symbol_statistics(\n    source_files=source_partitions,\n    median_features=median_features,\n    mode_features=mode_features,\n    imputed_dir=imputed_dir\n)\n\n# Calculate Global Fallbacks\nglobal_medians, global_modes = calculate_global_fallbacks(\n    symbol_medians=symbol_medians,\n    symbol_modes=symbol_modes,\n    median_features=median_features,\n    mode_features=mode_features\n)\n\n# Process and impute problematic partitions\nproblematic_results = process_and_impute_problematic_partitions(\n    partition_files=problematic_partition_files,\n    features=top_30_features,\n    output_dir=output_dir,\n    window=7,\n    symbol_medians=symbol_medians,\n    symbol_modes=symbol_modes,\n    global_medians=global_medians,\n    global_modes=global_modes\n)\n\n# Print summary of Phase 3 results\nprint(\"\\nPhase 3 Processing Results:\")\nfor res in problematic_results:\n    print(res)\n\nprint(\"\\n--- Phase 3 Completed ---\\n\")\n\n# Monitor memory usage\nprint_memory_usage()\n\n# Clean up\ndel problematic_results, symbol_medians, symbol_modes, global_medians, global_modes\ngc.collect()\n\n# ------------------------------ Phase 3 End ------------------------------\n","metadata":{"trusted":true,"execution":{"iopub.status.busy":"2024-12-30T07:46:19.845161Z","iopub.execute_input":"2024-12-30T07:46:19.846435Z","iopub.status.idle":"2024-12-30T07:46:51.169974Z","shell.execute_reply.started":"2024-12-30T07:46:19.846324Z","shell.execute_reply":"2024-12-30T07:46:51.168358Z"}},"outputs":[],"execution_count":null},{"cell_type":"markdown","source":"g. Verification and Merging","metadata":{}},{"cell_type":"code","source":"# ----------------------------- Verification Start -----------------------------\n\n# Define the directory containing the imputed partitions\noutput_dir = '/kaggle/working/imputed_partitions/'  # Ensure this matches the output_dir in the workflow\n\n# Define the problematic partitions and their respective features\nproblematic_partitions = {\n    'part.0.parquet': ['feature_26', 'feature_21', 'feature_27', 'feature_31', 'feature_04', 'feature_01'],\n    'part.1.parquet': ['feature_26', 'feature_21', 'feature_27', 'feature_31'],\n    'part.2.parquet': ['feature_26', 'feature_21', 'feature_27', 'feature_31']\n}\n\n# Define Top 30 Features (ensure consistency with the workflow)\ntop_30_features = [\n    'feature_61', 'feature_20', 'feature_24', 'feature_25', 'feature_23',\n    'feature_29', 'feature_22', 'feature_30', 'feature_38', 'feature_04',\n    'feature_05', 'feature_26', 'feature_28', 'feature_07', 'feature_01',\n    'feature_08', 'feature_21', 'feature_27', 'feature_31', 'feature_47',\n    'feature_58', 'feature_62', 'feature_37', 'feature_36', 'feature_60',\n    'feature_17', 'feature_64', 'feature_33', 'feature_49', 'feature_15'\n]\n\n# Define Helper Function for Verification\ndef verify_imputation_enhanced(output_dir, problematic_partitions, top_30_features):\n    \"\"\"\n    Verify that all nulls have been imputed in the specified partitions and features.\n\n    Parameters:\n    - output_dir (str): Directory containing imputed partitions.\n    - problematic_partitions (dict): Dictionary of problematic partitions and their features.\n    - top_30_features (list): List of top 30 feature names.\n\n    Returns:\n    - None\n    \"\"\"\n    for partition, features in problematic_partitions.items():\n        imputed_partition_file = os.path.basename(partition).replace('.parquet', '_imputed.parquet')\n        imputed_partition_path = os.path.join(output_dir, imputed_partition_file)\n        try:\n            ddf = dd.read_parquet(imputed_partition_path, columns=['symbol_id', 'time_id', 'date_id'] + features, engine='pyarrow')\n            df = ddf.compute()\n            print(f\"\\nVerifying {imputed_partition_file} with shape {df.shape}\")\n\n            # Check for missing columns\n            required_columns = ['symbol_id', 'time_id', 'date_id'] + top_30_features\n            missing_cols = set(required_columns) - set(df.columns)\n            if missing_cols:\n                print(f\"❌ Missing columns in {imputed_partition_file}: {missing_cols}\")\n            else:\n                print(f\"✅ All required columns are present in {imputed_partition_file}.\")\n\n            # Check for nulls in specified features\n            current_nulls = df[features].isnull().sum()\n            print(f\"Null counts in {imputed_partition_file}:\")\n            print(current_nulls)\n\n            # Check if all specified features have zero nulls\n            nulls_present = current_nulls[current_nulls > 0]\n            if nulls_present.empty:\n                print(f\"✅ No nulls remain in the specified features of {imputed_partition_file}.\")\n            else:\n                print(f\"❌ Nulls still present in the following features of {imputed_partition_file}:\")\n                print(nulls_present)\n\n            # Clean up\n            del df, ddf\n            gc.collect()\n\n        except Exception as e:\n            print(f\"Error verifying {imputed_partition_file}: {e}\")\n\n# Call the verification function\nverify_imputation_enhanced(output_dir, problematic_partitions, top_30_features)\n\n# ------------------------------ Verification End ------------------------------\n","metadata":{"trusted":true,"execution":{"iopub.status.busy":"2024-12-30T07:46:59.505925Z","iopub.execute_input":"2024-12-30T07:46:59.506393Z","iopub.status.idle":"2024-12-30T07:47:00.367747Z","shell.execute_reply.started":"2024-12-30T07:46:59.506357Z","shell.execute_reply":"2024-12-30T07:47:00.366290Z"}},"outputs":[],"execution_count":null},{"cell_type":"markdown","source":"g. Verification and Merging","metadata":{}},{"cell_type":"code","source":"# # ----------------------------- Verification Start -----------------------------\n\n# # Define the directory containing the imputed partitions\n# output_dir = '/kaggle/working/imputed_partitions/'  # Ensure this matches the output_dir in the workflow\n\n# # Define the problematic partitions and their respective features\n# problematic_partitions = {\n#     'part.0.parquet': ['feature_26', 'feature_21', 'feature_27', 'feature_31', 'feature_04', 'feature_01'],\n#     'part.1.parquet': ['feature_26', 'feature_21', 'feature_27', 'feature_31'],\n#     'part.2.parquet': ['feature_26', 'feature_21', 'feature_27', 'feature_31']\n# }\n\n# # Define Top 30 Features (ensure consistency with the workflow)\n# top_30_features = [\n#     'feature_61', 'feature_20', 'feature_24', 'feature_25', 'feature_23',\n#     'feature_29', 'feature_22', 'feature_30', 'feature_38', 'feature_04',\n#     'feature_05', 'feature_26', 'feature_28', 'feature_07', 'feature_01',\n#     'feature_08', 'feature_21', 'feature_27', 'feature_31', 'feature_47',\n#     'feature_58', 'feature_62', 'feature_37', 'feature_36', 'feature_60',\n#     'feature_17', 'feature_64', 'feature_33', 'feature_49', 'feature_15'\n# ]\n\n# # Define Helper Function for Verification\n# def verify_imputation_enhanced(output_dir, problematic_partitions, top_30_features):\n#     \"\"\"\n#     Verify that all nulls have been imputed in the specified partitions and features.\n\n#     Parameters:\n#     - output_dir (str): Directory containing imputed partitions.\n#     - problematic_partitions (dict): Dictionary of problematic partitions and their features.\n#     - top_30_features (list): List of top 30 feature names.\n\n#     Returns:\n#     - None\n#     \"\"\"\n#     for partition, features in problematic_partitions.items():\n#         imputed_partition_file = os.path.basename(partition).replace('.parquet', '_imputed.parquet')\n#         imputed_partition_path = os.path.join(output_dir, imputed_partition_file)\n#         try:\n#             ddf = dd.read_parquet(imputed_partition_path, columns=['symbol_id', 'time_id', 'date_id'] + features, engine='pyarrow')\n#             df = ddf.compute()\n#             print(f\"\\nVerifying {imputed_partition_file} with shape {df.shape}\")\n\n#             # Check for missing columns\n#             required_columns = ['symbol_id', 'time_id', 'date_id'] + top_30_features\n#             missing_cols = set(required_columns) - set(df.columns)\n#             if missing_cols:\n#                 print(f\"❌ Missing columns in {imputed_partition_file}: {missing_cols}\")\n#             else:\n#                 print(f\"✅ All required columns are present in {imputed_partition_file}.\")\n\n#             # Check for nulls in specified features\n#             current_nulls = df[features].isnull().sum()\n#             print(f\"Null counts in {imputed_partition_file}:\")\n#             print(current_nulls)\n\n#             # Check if all specified features have zero nulls\n#             nulls_present = current_nulls[current_nulls > 0]\n#             if nulls_present.empty:\n#                 print(f\"✅ No nulls remain in the specified features of {imputed_partition_file}.\")\n#             else:\n#                 print(f\"❌ Nulls still present in the following features of {imputed_partition_file}:\")\n#                 print(nulls_present)\n\n#             # Clean up\n#             del df, ddf\n#             gc.collect()\n\n#         except Exception as e:\n#             print(f\"Error verifying {imputed_partition_file}: {e}\")\n\n# # Call the verification function\n# verify_imputation_enhanced(output_dir, problematic_partitions, top_30_features)\n\n# # ------------------------------ Verification End ------------------------------\n","metadata":{"trusted":true},"outputs":[],"execution_count":null},{"cell_type":"markdown","source":"h. Step 3: Merge All Imputed Partitions","metadata":{}},{"cell_type":"code","source":"# ----------------------------- Step 3: Merge All Imputed Partitions -----------------------------\n\nimport dask.dataframe as dd\nimport os\nimport gc\n\n# Define the directory containing the imputed partitions\nimputed_dir = '/kaggle/working/imputed_partitions/'  # Ensure this is correct\n\n# Define the output directory for the merged DataFrame\nmerged_output_dir = '/kaggle/working/output_data/'\nos.makedirs(merged_output_dir, exist_ok=True)\n\n# Define the path pattern to read all imputed parquet files\nimputed_parquet_pattern = os.path.join(imputed_dir, 'part.*_imputed.parquet')\n\n# Read all imputed parquet files into a single Dask DataFrame\ntry:\n    merged_ddf = dd.read_parquet(imputed_parquet_pattern, engine='pyarrow')\n    print(f\"\\nSuccessfully read all imputed partitions into a single Dask DataFrame.\")\n    print(f\"Merged Dask DataFrame shape: {merged_ddf.shape}\")\nexcept Exception as e:\n    print(f\"Error reading imputed parquet files: {e}\")\n    merged_ddf = None\n\n# ------------------------------ Step 3 End ------------------------------\n","metadata":{"trusted":true,"execution":{"iopub.status.busy":"2024-12-30T07:48:55.327286Z","iopub.execute_input":"2024-12-30T07:48:55.327757Z","iopub.status.idle":"2024-12-30T07:48:55.389990Z","shell.execute_reply.started":"2024-12-30T07:48:55.327724Z","shell.execute_reply":"2024-12-30T07:48:55.388817Z"}},"outputs":[],"execution_count":null},{"cell_type":"markdown","source":"i. Step 4: Save Merged DataFrame","metadata":{}},{"cell_type":"code","source":"# ----------------------------- Step 4: Save Merged DataFrame -----------------------------\n\n# Define the path for the merged DataFrame\nmerged_parquet_path = os.path.join(merged_output_dir, 'merged_data.parquet')\nmerged_csv_path = os.path.join(merged_output_dir, 'merged_data.csv')\n\n# Save as Parquet using Dask\nif merged_ddf is not None:\n    try:\n        # Save as Parquet\n        merged_ddf.to_parquet(merged_parquet_path, compression='snappy', write_index=False)\n        print(f\"Successfully saved merged DataFrame to {merged_parquet_path}\")\n    except Exception as e:\n        print(f\"Error saving merged Parquet file: {e}\")\n\n    # # Optionally, save as CSV using Dask (Note: CSV saving is less efficient for large datasets)\n    # try:\n    #     merged_ddf.to_csv(os.path.join(merged_output_dir, 'merged_data.csv'), single_file=True, index=False)\n    #     print(f\"Successfully saved merged DataFrame to {merged_csv_path}\")\n    # except Exception as e:\n    #     print(f\"Error saving merged CSV file: {e}\")\nelse:\n    print(\"Merged Dask DataFrame is not available for saving.\")\n\n# Clean up\ndel merged_ddf\ngc.collect()\n\n# ------------------------------ Step 4 End ------------------------------\n","metadata":{"trusted":true,"execution":{"iopub.status.busy":"2024-12-30T07:50:35.907965Z","iopub.execute_input":"2024-12-30T07:50:35.908458Z","iopub.status.idle":"2024-12-30T07:51:35.069017Z","shell.execute_reply.started":"2024-12-30T07:50:35.908418Z","shell.execute_reply":"2024-12-30T07:51:35.067643Z"}},"outputs":[],"execution_count":null},{"cell_type":"code","source":"# ----------------------------- Verification After Merging -----------------------------\n\nimport dask.dataframe as dd\nimport os\n\n# Define the path for the merged DataFrame\nmerged_parquet_path = os.path.join(merged_output_dir, 'merged_data.parquet')\n\n# Load the merged DataFrame\ntry:\n    merged_ddf = dd.read_parquet(merged_parquet_path, engine='pyarrow')\n    print(f\"\\nLoaded merged DataFrame with shape: {merged_ddf.shape}\")\n    \n    # Compute the number of rows and columns\n    num_rows = merged_ddf.shape[0].compute()\n    num_cols = len(merged_ddf.columns)\n    print(f\"Merged DataFrame has {num_rows} rows and {num_cols} columns.\")\n    \n    # Optional: Preview the first few rows\n    print(\"\\nPreview of the merged DataFrame:\")\n    print(merged_ddf.head())\n    \nexcept Exception as e:\n    print(f\"Error loading merged DataFrame: {e}\")\n\n# Clean up\ndel merged_ddf\ngc.collect()\n\n# ------------------------------ Verification End ------------------------------\n","metadata":{"trusted":true,"execution":{"iopub.status.busy":"2024-12-30T07:52:55.487179Z","iopub.execute_input":"2024-12-30T07:52:55.487678Z","iopub.status.idle":"2024-12-30T07:52:56.891994Z","shell.execute_reply.started":"2024-12-30T07:52:55.487642Z","shell.execute_reply":"2024-12-30T07:52:56.890485Z"}},"outputs":[],"execution_count":null},{"cell_type":"code","source":"","metadata":{"trusted":true},"outputs":[],"execution_count":null},{"cell_type":"code","source":"","metadata":{"trusted":true},"outputs":[],"execution_count":null},{"cell_type":"code","source":"","metadata":{"trusted":true},"outputs":[],"execution_count":null},{"cell_type":"code","source":"","metadata":{"trusted":true},"outputs":[],"execution_count":null},{"cell_type":"code","source":"","metadata":{"trusted":true},"outputs":[],"execution_count":null},{"cell_type":"code","source":"","metadata":{"trusted":true},"outputs":[],"execution_count":null}]}