{"metadata":{"kernelspec":{"language":"python","display_name":"Python 3","name":"python3"},"language_info":{"name":"python","version":"3.10.14","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":30786,"isInternetEnabled":true,"language":"python","sourceType":"notebook","isGpuEnabled":false}},"nbformat_minor":4,"nbformat":4,"cells":[{"cell_type":"code","source":"import os, gc, joblib\nimport numpy as np\nimport pandas as pd\nfrom tqdm.notebook import tqdm\nfrom pathlib import Path","metadata":{"_uuid":"8f2839f25d086af736a60e9eeb907d3b93b6e0e5","_cell_guid":"b1076dfc-b9ad-4769-8c92-a6c4dae69d19","execution":{"iopub.status.busy":"2024-10-23T13:28:16.307523Z","iopub.execute_input":"2024-10-23T13:28:16.307931Z","iopub.status.idle":"2024-10-23T13:28:20.880442Z","shell.execute_reply.started":"2024-10-23T13:28:16.307861Z","shell.execute_reply":"2024-10-23T13:28:20.879258Z"},"trusted":true},"execution_count":null,"outputs":[]},{"cell_type":"markdown","source":"# Data Preparation Notebook\n\nThis notebook makes a folds Kaggle Dataset that you can use downstream in training. It's helpful to process the data like this if you are compute contrained (i.e., only using Kaggle notebooks).\n\nLoad each partition; check summary stats; combine pairs of partitions as folds to use in training notebook.","metadata":{}},{"cell_type":"code","source":"def reduce_mem_usage(df, verbose=True):\n    start_mem = df.memory_usage().sum() / 1024**2\n    if verbose:\n        print(f\"Memory usage of dataframe is {start_mem:.2f} MB\")\n    \n    for col in df.columns:\n        col_type = df[col].dtype\n        \n        if col_type != object:  # Exclude strings\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                elif c_min > np.iinfo(np.int64).min and c_max < np.iinfo(np.int64).max:\n                    df[col] = df[col].astype(np.int64)\n            \n            elif str(col_type)[:5] == 'float':\n                if c_min > np.finfo(np.float16).min and c_max < np.finfo(np.float16).max:\n                    df[col] = df[col].astype(np.float16)\n                elif c_min > np.finfo(np.float32).min and c_max < np.finfo(np.float32).max:\n                    df[col] = df[col].astype(np.float32)\n                else:\n                    df[col] = df[col].astype(np.float64)\n        else:\n            df[col] = df[col].astype('category')\n\n    end_mem = df.memory_usage().sum() / 1024**2\n    if verbose:\n        print(f\"Memory usage after optimization is {end_mem:.2f} MB\")\n        print(f\"Decreased by {100 * (start_mem - end_mem) / start_mem:.1f}%\")\n    \n    return df","metadata":{"execution":{"iopub.status.busy":"2024-10-23T13:28:20.882445Z","iopub.execute_input":"2024-10-23T13:28:20.883372Z","iopub.status.idle":"2024-10-23T13:28:20.904000Z","shell.execute_reply.started":"2024-10-23T13:28:20.883311Z","shell.execute_reply":"2024-10-23T13:28:20.902073Z"},"trusted":true},"execution_count":null,"outputs":[]},{"cell_type":"code","source":"def get_partition_summary(input_path, partitions=range(10)):\n    # Initialize an empty list to store the partition summaries\n    summary_list = []\n\n    for partition in tqdm(partitions, desc=\"Processing Partitions\"):\n        # Construct the path for the current partition (assumes partition directories follow a specific pattern)\n        partition_path = input_path / f'train.parquet/partition_id={partition}'\n        \n        # Read the partition\n        df = pd.read_parquet(partition_path)\n        \n        # Calculate the required stats for each partition\n        row_count = len(df)\n        start_date = df['date_id'].min()\n        end_date = df['date_id'].max()\n        memory_usage = df.memory_usage().sum() / 1024**2  # Memory usage in MB\n        df = reduce_mem_usage(df, verbose=False)\n        memory_usage_reduced = df.memory_usage().sum() / 1024**2  # Memory usage in MB\n        \n        # Calculate percentage of NaN values\n        total_elements = df.size  # Total number of elements in the DataFrame (rows * columns)\n        nan_count = df.isna().sum().sum()  # Total number of NaN values\n        nan_percentage = (nan_count / total_elements) * 100  # Percentage of NaN values\n        \n        # Append the result for this partition\n        summary_list.append({\n            'Partition': partition,\n            'Row Count': row_count,\n            'Start Date': start_date,\n            'End Date': end_date,\n            'Days': end_date - start_date,\n            'Num Symbols': df['symbol_id'].nunique(),\n            '% NaN': nan_percentage,\n            'Memory Usage (MB)': memory_usage,\n            'Memory Reduced (MB)': memory_usage_reduced\n        })\n        \n        del df\n        gc.collect()\n    \n    # Convert the list of dictionaries into a DataFrame for easy display\n    summary_df = pd.DataFrame(summary_list)\n    \n    # Create a DataFrame for totals\n    totals_df = pd.DataFrame([{\n        'Partition': 'Total',\n        'Row Count': summary_df['Row Count'].sum(),\n        'Start Date': '',  # No total for date columns\n        'End Date': '',    # No total for date columns\n        'Days': summary_df['End Date'].max() - summary_df['Start Date'].min(),\n        'Num Symbols': '',\n        'Memory Usage (MB)': summary_df['Memory Usage (MB)'].sum(),\n        'Memory Reduced (MB)': summary_df['Memory Reduced (MB)'].sum(),\n        '% NaN': (summary_df['% NaN'] * summary_df['Row Count']).sum() / summary_df['Row Count'].sum()  # Weighted average of NaN percentage\n    }])\n\n    # Concatenate the totals row with the original summary dataframe\n    summary_df = pd.concat([summary_df, totals_df], ignore_index=True)\n    \n    return summary_df","metadata":{"execution":{"iopub.status.busy":"2024-10-23T13:28:20.907382Z","iopub.execute_input":"2024-10-23T13:28:20.907957Z","iopub.status.idle":"2024-10-23T13:28:20.925420Z","shell.execute_reply.started":"2024-10-23T13:28:20.907897Z","shell.execute_reply":"2024-10-23T13:28:20.923967Z"},"trusted":true},"execution_count":null,"outputs":[]},{"cell_type":"markdown","source":"We cannot load all the data because we get an OOM (try it if you don't belive me). Let's look at each partition info.","metadata":{}},{"cell_type":"code","source":"input_path = Path('/kaggle/input/jane-street-real-time-market-data-forecasting')","metadata":{"execution":{"iopub.status.busy":"2024-10-23T13:28:20.926983Z","iopub.execute_input":"2024-10-23T13:28:20.927492Z","iopub.status.idle":"2024-10-23T13:28:20.941076Z","shell.execute_reply.started":"2024-10-23T13:28:20.927424Z","shell.execute_reply":"2024-10-23T13:28:20.939482Z"},"trusted":true},"execution_count":null,"outputs":[]},{"cell_type":"code","source":"%%time\n\nsummary_df = get_partition_summary(input_path)","metadata":{"execution":{"iopub.status.busy":"2024-10-23T13:28:20.945644Z","iopub.execute_input":"2024-10-23T13:28:20.946147Z","iopub.status.idle":"2024-10-23T13:30:05.818460Z","shell.execute_reply.started":"2024-10-23T13:28:20.946096Z","shell.execute_reply":"2024-10-23T13:30:05.817284Z"},"trusted":true},"execution_count":null,"outputs":[]},{"cell_type":"code","source":"summary_df","metadata":{"execution":{"iopub.status.busy":"2024-10-23T13:30:05.820108Z","iopub.execute_input":"2024-10-23T13:30:05.821203Z","iopub.status.idle":"2024-10-23T13:30:05.849181Z","shell.execute_reply.started":"2024-10-23T13:30:05.821146Z","shell.execute_reply":"2024-10-23T13:30:05.847798Z"},"trusted":true},"execution_count":null,"outputs":[]},{"cell_type":"markdown","source":"From the above, we see that the full dataset is ~16GB. We can reduce that to ~8GB. We _should_ be able to load full data, but my experience is that pandas makes copies and is very inefficient. Loading the full (reduced) dataset and then doing feature engineering and model fitting will probably OOM. The first few partitions have big percent of NaNs; perhaps the production system or something evolved and these early partitions are not representative of how things are today. \n\nOne strategy is to make folds simply by concatenating partitions. That's what the folllowing code does, however it is disabled in this run.","metadata":{}},{"cell_type":"code","source":"def read_and_concatenate_partitions(input_path, output_path, partition_pairs):\n    # Create the output directory if it doesn't exist\n    output_path.mkdir(parents=True, exist_ok=True)\n    \n    # Loop through each pair of partitions\n    for pair in tqdm(partition_pairs):\n        partition_1, partition_2 = pair\n\n        # Read both partitions\n        partition_1_path = input_path / f'train.parquet/partition_id={partition_1}'\n        partition_2_path = input_path / f'train.parquet/partition_id={partition_2}'\n\n        df1 = pd.read_parquet(partition_1_path)\n        df2 = pd.read_parquet(partition_2_path)\n        \n        # Concatenate the two DataFrames\n        concatenated_df = pd.concat([df1, df2], ignore_index=True)\n\n        # Reduce memory usage for both partitions\n        concatenated_df = reduce_mem_usage(concatenated_df, verbose=False)\n\n        # Write the concatenated DataFrame to a new Parquet partition\n        output_partition_path = output_path / f'partitions_{partition_1}_{partition_2}.parquet'\n        concatenated_df.to_parquet(output_partition_path, index=False)\n        \n        del df1\n        del df2\n        del concatenated_df\n        \n        gc.collect()\n\n        print(f\"Written concatenated partition {partition_1} and {partition_2} to {output_partition_path}\")","metadata":{"execution":{"iopub.status.busy":"2024-10-23T13:30:05.851724Z","iopub.execute_input":"2024-10-23T13:30:05.852485Z","iopub.status.idle":"2024-10-23T13:30:05.862898Z","shell.execute_reply.started":"2024-10-23T13:30:05.852436Z","shell.execute_reply":"2024-10-23T13:30:05.860772Z"},"trusted":true},"execution_count":null,"outputs":[]},{"cell_type":"code","source":"# Define input and output paths\noutput_path = Path('/kaggle/working/folds')\n\n# List of partition pairs to process\npartition_pairs = [(2, 3), (4, 5), (6, 7), (8, 9)]\n\n# Process and write new Parquet partitions \n# read_and_concatenate_partitions(input_path, output_path, partition_pairs)","metadata":{"execution":{"iopub.status.busy":"2024-10-23T13:30:05.864805Z","iopub.execute_input":"2024-10-23T13:30:05.865284Z","iopub.status.idle":"2024-10-23T13:30:05.881141Z","shell.execute_reply.started":"2024-10-23T13:30:05.865239Z","shell.execute_reply":"2024-10-23T13:30:05.879600Z"},"trusted":true},"execution_count":null,"outputs":[]},{"cell_type":"markdown","source":"# Make Single Dataset\n\nHere I make a single dataset, memory reduced. I drop the first three partitions because they seem to be materially different than the rest.","metadata":{}},{"cell_type":"code","source":"partitions_list = [3, 4, 5, 6, 7, 8, 9]\noutput_path = Path('/kaggle/working/')","metadata":{"execution":{"iopub.status.busy":"2024-10-23T13:30:06.430749Z","iopub.execute_input":"2024-10-23T13:30:06.431235Z","iopub.status.idle":"2024-10-23T13:30:06.436915Z","shell.execute_reply.started":"2024-10-23T13:30:06.431189Z","shell.execute_reply":"2024-10-23T13:30:06.435597Z"},"trusted":true},"execution_count":null,"outputs":[]},{"cell_type":"code","source":"def read_and_concatenate_all_partitions(input_path, output_path, partitions_list):\n    # Create the output directory if it doesn't exist\n    output_path.mkdir(parents=True, exist_ok=True)\n    \n    partitions = []\n    \n    # Loop through each pair of partitions\n    for partition in tqdm(partitions_list):\n\n        partition_path = input_path / f'train.parquet/partition_id={partition}'\n\n        df = pd.read_parquet(partition_path)\n        df = reduce_mem_usage(df, verbose=False)\n\n        partitions.append(df)\n        \n        del df\n        gc.collect()\n        \n    all_dfs = pd.concat(partitions, ignore_index=True)\n    output_partition_path = output_path / 'partitions_reduced.parquet'\n    all_dfs.to_parquet(output_partition_path, index=False)\n    \n    del all_dfs\n    gc.collect()\n    ","metadata":{"execution":{"iopub.status.busy":"2024-10-23T13:30:07.912712Z","iopub.execute_input":"2024-10-23T13:30:07.913204Z","iopub.status.idle":"2024-10-23T13:30:07.922585Z","shell.execute_reply.started":"2024-10-23T13:30:07.913158Z","shell.execute_reply":"2024-10-23T13:30:07.921168Z"},"trusted":true},"execution_count":null,"outputs":[]},{"cell_type":"code","source":"%%time\n\nread_and_concatenate_all_partitions(input_path, output_path, partitions_list)","metadata":{"execution":{"iopub.status.busy":"2024-10-22T21:13:59.087939Z","iopub.execute_input":"2024-10-22T21:13:59.088447Z","iopub.status.idle":"2024-10-22T21:19:53.368666Z","shell.execute_reply.started":"2024-10-22T21:13:59.088392Z","shell.execute_reply":"2024-10-22T21:19:53.367057Z"},"trusted":true},"execution_count":null,"outputs":[]},{"cell_type":"markdown","source":"# Make Full Period Equal Folds\n\nThis fold strategy follows the comment thread https://www.kaggle.com/competitions/jane-street-real-time-market-data-forecasting/discussion/542057\n\nNote that instead of explicitly skipping the first 500 days, we are just excluding the first three partitions (510 days).","metadata":{}},{"cell_type":"code","source":"valid_days = 100\nn_folds = 5","metadata":{"execution":{"iopub.status.busy":"2024-10-23T13:31:09.259374Z","iopub.execute_input":"2024-10-23T13:31:09.259858Z","iopub.status.idle":"2024-10-23T13:31:09.265119Z","shell.execute_reply.started":"2024-10-23T13:31:09.259807Z","shell.execute_reply":"2024-10-23T13:31:09.263812Z"},"trusted":true},"execution_count":null,"outputs":[]},{"cell_type":"code","source":"reduced_partition = pd.read_parquet(output_path / 'partitions_reduced.parquet')\n\nreduced_partition.info()","metadata":{"execution":{"iopub.status.busy":"2024-10-22T21:19:53.378626Z","iopub.execute_input":"2024-10-22T21:19:53.379055Z","iopub.status.idle":"2024-10-22T21:20:37.938376Z","shell.execute_reply.started":"2024-10-22T21:19:53.379013Z","shell.execute_reply":"2024-10-22T21:20:37.937240Z"},"trusted":true},"execution_count":null,"outputs":[]},{"cell_type":"code","source":"output_path = Path('/kaggle/working/folds')\n\ndates = reduced_partition['date_id'].unique()\n\n# Definalle validation dates as the last `num_valid_dates` dates\nvalid_dates = dates[-valid_days:]\n\n# Define training dates as all dates except the last `num_valid_dates` dates\ntrain_dates = dates[:-valid_days]","metadata":{},"execution_count":null,"outputs":[]},{"cell_type":"code","source":"def make_folds(n_folds, df):\n    output_path.mkdir(parents=True, exist_ok=True)\n    for i in tqdm(range(n_folds)):\n        selected_dates = [date for ii, date in enumerate(train_dates) if ii % n_folds == i]\n        fold_df = df.loc[df['date_id'].isin(selected_dates)]\n        \n        print(f\"fold {i}: min_date: {fold_df['date_id'].min()}\\t max_date: {fold_df['date_id'].max()}\\t n_rows: {len(fold_df)}\")\n        fold_df.to_parquet(output_path / f'fold_{i}.parquet', index=False)\n        del fold_df\n        gc.collect()\n    \n    valid_df = df.loc[df['date_id'].isin(valid_dates)]\n    valid_df.to_parquet(output_path / f'valid.parquet', index=False)","metadata":{},"execution_count":null,"outputs":[]},{"cell_type":"code","source":"%%time\n\nmake_folds(n_folds, reduced_partition)","metadata":{},"execution_count":null,"outputs":[]},{"cell_type":"markdown","source":"Good luck, and may all your overfitting be benign.","metadata":{}}]}