{"metadata":{"kernelspec":{"language":"python","display_name":"Python 3","name":"python3"},"language_info":{"pygments_lexer":"ipython3","nbconvert_exporter":"python","version":"3.6.4","file_extension":".py","codemirror_mode":{"name":"ipython","version":3},"name":"python","mimetype":"text/x-python"}},"nbformat_minor":4,"nbformat":4,"cells":[{"cell_type":"code","source":"!pip install /kaggle/input/polars01516/polars-0.15.16-cp37-abi3-manylinux_2_17_x86_64.manylinux2014_x86_64.whl\n# Move software to working disk\n!rm  -r software\n!scp -r /kaggle/input/graphnet-and-dependencies/software .\n\n# Install dependencies\n!pip install /kaggle/working/software/dependencies/torch-1.11.0+cu115-cp37-cp37m-linux_x86_64.whl\n!pip install /kaggle/working/software/dependencies/torch_cluster-1.6.0-cp37-cp37m-linux_x86_64.whl\n!pip install /kaggle/working/software/dependencies/torch_scatter-2.0.9-cp37-cp37m-linux_x86_64.whl\n!pip install /kaggle/working/software/dependencies/torch_sparse-0.6.13-cp37-cp37m-linux_x86_64.whl\n!pip install /kaggle/working/software/dependencies/torch_geometric-2.0.4.tar.gz\n\n# Install GraphNeT\n!cd software/graphnet;pip install --no-index --find-links=\"/kaggle/working/software/dependencies\" -e .[torch]\n\n# Append to PATH\nimport sys\nsys.path.append('/kaggle/working/software/graphnet/src')","metadata":{"_uuid":"8f2839f25d086af736a60e9eeb907d3b93b6e0e5","_cell_guid":"b1076dfc-b9ad-4769-8c92-a6c4dae69d19","scrolled":true,"execution":{"iopub.status.busy":"2023-02-14T15:18:11.620188Z","iopub.execute_input":"2023-02-14T15:18:11.620656Z","iopub.status.idle":"2023-02-14T15:23:39.617728Z","shell.execute_reply.started":"2023-02-14T15:18:11.620617Z","shell.execute_reply":"2023-02-14T15:23:39.615820Z"},"trusted":true},"execution_count":null,"outputs":[]},{"cell_type":"code","source":"## Configuration parameters\nMODE = 'train'\n# MODE = 'test'\n\n# USE_POLARS = False\nUSE_POLARS = True\n\n# TRAIN_MAX_EVENTS = 20000\nTRAIN_MAX_EVENTS = None\nTRAIN_BATCH_START = 1\n# TRAIN_N_BATCHES = 1\n\n## I pulled in one piece of older code to demonstrate \"before and after\".\n## Set to True (and USE_POLARS=False) if interested in seeing the difference.\nUSE_UNOPTIMIZED = False\n\n\n#### HYPERPARAMETERS ####\n\n## For setting auxiliary = False\n## Hand-tuned and hand-validated, I mostly used batches 100-105, and probably early on also touched batch 1.\n## TODO: revisit the deep core logic now that I've learned about the deep veto layer: https://www.kaggle.com/competitions/icecube-neutrinos-in-deep-ice/discussion/381702\nFIND_BEST_POINTS = True\nMIN_PRIMARY_DATAPOINTS = 2\nMAX_Z = 3\nMAX_DEEP_Z = 1\nMAX_T = 350\nMAX_DEEP_T = 180\nif USE_POLARS:\n    ## Only implmented in polars version.\n    AUX_FALSE_WEIGHT = 0.01\n\n## Ensemble and algorithm selection parameters.\n## First weight is center of charge algorithm introduced in this notebook.\n## Second weight is the unweighted version. Which turns out to be the least-squares algorithm: https://www.kaggle.com/competitions/icecube-neutrinos-in-deep-ice/discussion/381747\nUSE_ENSEMBLE = True\nWEIGHTS = [0.575, 0.425]\nif not USE_ENSEMBLE and not USE_POLARS:\n    ALGORITHM = 'center of charge'\n#     ALGORITHM = 'least squares'\n    if ALGORITHM == 'least squares':\n        USE_WEIGHTED_LEAST_SQUARES = False\n\n\n## Constants\nINPUT_DIR = '/kaggle/input/icecube-neutrinos-in-deep-ice'\n\n## Basic configuration override logic\nif MODE == 'test':\n    TRAIN_MAX_EVENTS = None\nif USE_ENSEMBLE:\n    USE_WEIGHTED_LEAST_SQUARES = False","metadata":{"execution":{"iopub.status.busy":"2023-02-14T15:23:39.622596Z","iopub.execute_input":"2023-02-14T15:23:39.623104Z","iopub.status.idle":"2023-02-14T15:23:39.635111Z","shell.execute_reply.started":"2023-02-14T15:23:39.623044Z","shell.execute_reply":"2023-02-14T15:23:39.633562Z"},"trusted":true},"execution_count":null,"outputs":[]},{"cell_type":"code","source":"import graphnet\nimport pyarrow.parquet as pq\nimport sqlite3\nimport pandas as pd\nimport sqlalchemy\nfrom tqdm import tqdm\nimport os\nfrom typing import Any, Dict, List, Optional\nimport polars as pl\nimport math\nimport time\nimport gc\nimport numpy as np\ntqdm.pandas()\ntry:\n    from graphnet.data.sqlite.sqlite_utilities import create_table\nexcept:\n    from graphnet.data.sqlite.sqlite_utilities import create_table \n\ndef load_input(meta_batch: pl.DataFrame, input_data_folder: str) -> pl.DataFrame:\n        \"\"\"\n        Will load the corresponding detector readings associated with the meta data batch.\n        \"\"\"\n        batch_id = meta_batch.select(pl.col('batch_id')).unique()\n\n        assert batch_id.height == 1, \"contains multiple batch_ids. Did you set the batch_size correctly?\"\n\n\n        detector_readings = batch_integration(batch_id.item())\n\n        \n#         sensor_positions = geometry_table.collect()[detector_readings['sensor_id'], ['x', 'y', 'z']]\n#         sensor_positions.index = detector_readings.index\n\n#         for column in sensor_positions.columns:\n#             if column not in detector_readings.columns:\n#                 detector_readings = detector_readings.with_columns(sensor_positions[column])\n\n#         detector_readings = detector_readings.with_columns(\n#             pl.when(pl.col(\"auxiliary\") == False).then(0).when(\n#                 pl.col(\"auxiliary\") == True).then(1).alias('auxiliary'))\n#         detector_readings\n#         print(detector_readings.collect())\n        return detector_readings\n","metadata":{"execution":{"iopub.status.busy":"2023-02-14T15:23:39.637189Z","iopub.execute_input":"2023-02-14T15:23:39.637711Z","iopub.status.idle":"2023-02-14T15:23:41.040706Z","shell.execute_reply.started":"2023-02-14T15:23:39.637660Z","shell.execute_reply":"2023-02-14T15:23:41.039356Z"},"trusted":true},"execution_count":null,"outputs":[]},{"cell_type":"code","source":"def least_squares(batch, weighted=False):\n    batch['xt'] = batch.x * batch.time\n    batch['yt'] = batch.y * batch.time\n    batch['zt'] = batch.z * batch.time\n    batch['tt'] = batch.time * batch.time\n    if weighted:\n        df = batch[['x', 'y', 'z', 'time', 'xt', 'yt', 'zt', 'tt']] * batch.charge.values[:, None]\n        df['charge'] = batch.charge\n        df = df.groupby('event_id').sum()\n        df = df.div(df.charge, axis=0)\n    else:\n        df = batch[['x', 'y', 'z', 'time', 'xt', 'yt', 'zt', 'tt']]\n        df = df.groupby('event_id').mean()\n    df[['x', 'y', 'z']] = (\n                              (df[['xt', 'yt', 'zt']].values - (df[['x', 'y', 'z']].values * df['time'].values[:, None]))\n                            / (df['tt'].values - (df.time.values * df.time.values))[:, None]\n                          )\n    ## Reverse it\n    df = -df[['x', 'y', 'z']]\n    df[['azimuth', 'zenith']] = angles_from_vectors(df.values)\n    return df[['azimuth', 'zenith']]\n","metadata":{"execution":{"iopub.status.busy":"2023-02-14T15:23:41.043189Z","iopub.execute_input":"2023-02-14T15:23:41.043736Z","iopub.status.idle":"2023-02-14T15:23:41.056341Z","shell.execute_reply.started":"2023-02-14T15:23:41.043678Z","shell.execute_reply":"2023-02-14T15:23:41.054453Z"},"trusted":true},"execution_count":null,"outputs":[]},{"cell_type":"code","source":"def process_batch(batch_id, sensor, max_events=None):\n    print('load batch...')\n    batch = pd.read_parquet(f'{INPUT_DIR}/{MODE}/batch_{batch_id}.parquet')\n\n    ## Limit to max_events\n    if max_events is not None:\n        batch_i = batch.reset_index()\n        event_ids = batch_i.event_id.drop_duplicates()\n        end_index = event_ids.index[max_events]\n        batch = batch_i[:end_index].set_index('event_id')\n    print(batch.shape)\n\n    ## Merge in sensor x,y,z data\n    batch = batch.reset_index().merge(sensor, how='left', on='sensor_id', left_index=False).set_index('event_id')\n\n    ## The logic for auxiliary = False is very basic, we can improve on it. Initial discussion here: \n    if FIND_BEST_POINTS:\n        batch = find_best_points_pandas(batch)\n \n    ## Limit to primary (aux=False) datapoints if there's enough of them.\n    ## For event_ids with too few, set all rows to aux=False. This handles these cases without any for loop logic.\n    ## MIN_PRIMARY_DATAPOINTS is a tuned value, based on optimizing the score on batches 101-110.\n    ## But the result was 'as low as possible'.\n    batch['primary_count'] = batch.groupby('event_id')['auxiliary'].transform('count') - batch.groupby('event_id')['auxiliary'].transform('sum')\n    batch.loc[batch.primary_count < MIN_PRIMARY_DATAPOINTS, 'auxiliary'] = False\n    batch = batch[batch.auxiliary == False]\n    print(batch.shape)\n\n    if USE_ENSEMBLE or ALGORITHM == 'center of charge':\n        df = center_of_charge(batch)\n        if USE_ENSEMBLE:\n            df1 = df\n    if USE_ENSEMBLE or ALGORITHM == 'least squares':\n        df = least_squares(batch, weighted=USE_WEIGHTED_LEAST_SQUARES)\n    if USE_ENSEMBLE:\n        df[['azimuth', 'zenith']] = average_angles([df1.azimuth.values, df.azimuth.values],\n                                                   [df1.zenith.values, df.zenith.values], \n                                                   weights=WEIGHTS)\n    return df","metadata":{"execution":{"iopub.status.busy":"2023-02-14T15:23:41.061040Z","iopub.execute_input":"2023-02-14T15:23:41.061527Z","iopub.status.idle":"2023-02-14T15:23:41.081338Z","shell.execute_reply.started":"2023-02-14T15:23:41.061487Z","shell.execute_reply":"2023-02-14T15:23:41.080245Z"},"trusted":true},"execution_count":null,"outputs":[]},{"cell_type":"code","source":"def proximity(e, col):\n    ## Since we explode the data to an NxN array, if N is too high it will take a long time and anyways we'll run out of memory.\n    ## TODO: see how high we can go before we run out of memory\n    ## TODO: consider subsampling instead of simply returning the baseline auxiliary setting?\n    ##       And/or splitting on string_id to get much smaller groups\n    if e.shape[0] > 2000:\n        return e[:, col['auxiliary']]\n\n    ## The magic None in the line below is a new thing I learned while working on this notebook.\n    ## It is shorthand for np.newaxis, and broadcasts the array to another dimension.\n    ## This allows us to create an NxN array for each event, so that we can\n    ## check if *any* other row in the event meets our proximity in time and height requirements.\n    ## For any matching pairs, then we set both rows as auxiliary = False. Rows without matches are auxiliary = True.\n    deltas = np.abs(e[:, [col['string_id'], col['depth_id'], col['time']]] - e[:, None, [col['string_id'], col['depth_id'], col['time']]])\n    dz = deltas[:, :, 1]\n\n    ## if same depth or different string id, ignore by setting dz > the max threshold used later.\n    dz[(dz == 0) | (deltas[:, :, 0] != 0)] = MAX_Z + MAX_DEEP_Z + 1\n\n    ## if sensor is not a deep ice sensor, and time > MAX_T, ignore\n    mask = (e[:, col['sensor_id']] < 4680)\n    mask = np.broadcast_to(mask, (mask.shape[0], mask.shape[0])).T\n    dz[mask & (deltas[:, :, 2] > MAX_T)] = MAX_Z + MAX_DEEP_Z + 1\n    ## if sensor IS a deep ice sensor, and time > MAX_DEEP_T, ignore\n    mask = (e[:, col['sensor_id']] >= 4680)\n    mask = np.broadcast_to(mask, (mask.shape[0], mask.shape[0])).T\n    dz[mask & (deltas[:, :, 2] > MAX_DEEP_T)] = MAX_Z + MAX_DEEP_Z + 1\n\n    ## Now take the min (best) result for each row com\n    dz = dz.min(axis=1)\n    ## If no matches, the default, then everything is aux=True\n    e[:, col['auxiliary']] = True\n    ## If not deep ice and distance less than threshold, or deep ice and distance less than other threshold, then we have a match!\n    e[((e[:, col['sensor_id']] < 4680) & (dz <= MAX_Z)) | ((e[:, col['sensor_id']] >= 4680) & (dz <= MAX_DEEP_Z)), col['auxiliary']] = False\n    ## Return only the data needed to speed up the np.concatenate called next.\n    return e[:, col['auxiliary']]","metadata":{"execution":{"iopub.status.busy":"2023-02-14T15:23:41.083065Z","iopub.execute_input":"2023-02-14T15:23:41.083563Z","iopub.status.idle":"2023-02-14T15:23:41.100201Z","shell.execute_reply.started":"2023-02-14T15:23:41.083517Z","shell.execute_reply":"2023-02-14T15:23:41.099160Z"},"trusted":true},"execution_count":null,"outputs":[]},{"cell_type":"code","source":"## We use the same numpy proximity function, the only difference is we convert from and to a Polars df instead of a Pandas df.\ndef find_best_points_polars(batch):\n    ## Minimize the size of the inputs, and make sure all data types are the same for an efficient np array\n    df = batch.select([pl.col('sensor_id'), pl.col('time'), pl.col('auxiliary'), pl.col('string_id'), pl.col('depth_id')])\n    column_to_index = { k:v for v,k in enumerate(df.columns)}\n    batch_values = df.to_numpy().astype('float32')\n\n    ## Pure numpy equivalent of groupby\n    events = np.split(batch_values, np.unique(batch.select(pl.col('event_id')), return_index=True)[1][1:])\n\n    ## Run each event sequentially in a list comprehension, then join back together with np.concatenate. So far tried and failed to find a reasonable solution to avoid this groupby > apply > join loop.\n    ## Rename 'auxiliary' -> 'best', and flip the boolean value\n    batch = batch.with_columns(pl.Series(np.concatenate([proximity(e, column_to_index) for e in tqdm(events)]).astype('bool')).alias('best')).with_columns(pl.col('best').is_not())\n\n    return batch","metadata":{"execution":{"iopub.status.busy":"2023-02-14T15:23:41.101822Z","iopub.execute_input":"2023-02-14T15:23:41.102298Z","iopub.status.idle":"2023-02-14T15:23:41.117189Z","shell.execute_reply.started":"2023-02-14T15:23:41.102179Z","shell.execute_reply":"2023-02-14T15:23:41.116279Z"},"trusted":true},"execution_count":null,"outputs":[]},{"cell_type":"code","source":"def time_weighted_centering(batch: pl, charge_weighted=True):\n    ## Polars equivalent to groupby->transform is called 'over'. We again use this to get min and max without a for loop.\n    batch = batch.with_columns([pl.col('time').min().over('event_id').alias('ev_t_min'),\n                                pl.col('time').max().over('event_id').alias('ev_t_max')])\n    if charge_weighted:\n        batch = batch.with_columns((pl.col('charge') * (pl.col('time') - pl.col('ev_t_min'))\n                                    / (pl.col('ev_t_max') - pl.col('ev_t_min'))).alias('w1'))\n        batch = batch.with_columns((pl.col('charge') - pl.col('w1')).alias('w0'))\n    else:\n        batch = batch.with_columns(((pl.col('time') - pl.col('ev_t_min'))\n                                    / (pl.col('ev_t_max') - pl.col('ev_t_min'))).alias('w1'))\n        batch = batch.with_columns((pl.lit(1) - pl.col('w1')).alias('w0'))\n\n    batch = batch.select(\n        [\n            pl.col('event_id'),\n            pl.col('w0'),\n            pl.col('w1'),\n            (pl.col('x') * pl.col('w0')).alias('wx0'),\n            (pl.col('y') * pl.col('w0')).alias('wy0'),\n            (pl.col('z') * pl.col('w0')).alias('wz0'),\n            (pl.col('x') * pl.col('w1')).alias('wx1'),\n            (pl.col('y') * pl.col('w1')).alias('wy1'),\n            (pl.col('z') * pl.col('w1')).alias('wz1'),\n        ]\n    ).collect().groupby('event_id', maintain_order=True).sum()\n\n    ## The direction the neutrino is traveling FROM is point0 - point1, instead of point1 - point0.\n    ## Counter-intuitive to me, but fortunately, easy to notice and correct if your score is > 1.57 instead of less.\n    batch = batch.with_columns([\n    ((pl.col('wx0') / pl.col('w0')) -\n     (pl.col('wx1') / pl.col('w1'))).alias('center_x'),\n    ((pl.col('wy0') / pl.col('w0')) -\n     (pl.col('wy1') / pl.col('w1'))).alias('center_y'),\n    ((pl.col('wz0') / pl.col('w0')) -\n     (pl.col('wz1') / pl.col('w1'))).alias('center_z'),\n    ])\n    \n    return batch\n","metadata":{"execution":{"iopub.status.busy":"2023-02-14T15:23:41.119314Z","iopub.execute_input":"2023-02-14T15:23:41.120208Z","iopub.status.idle":"2023-02-14T15:23:41.136677Z","shell.execute_reply.started":"2023-02-14T15:23:41.120162Z","shell.execute_reply":"2023-02-14T15:23:41.135351Z"},"trusted":true},"execution_count":null,"outputs":[]},{"cell_type":"code","source":"def batch_integration(batch_id):\n    print(MODE)\n    ## Scan parquet is part of Polars lazy evaluation, so no cost yet,\n    ## and it can figure out optimizations when we finally 'collect' it later.\n    meta = pl.scan_parquet(f'{INPUT_DIR}/{MODE}_meta.parquet')\n\n    print('load sensor data...')\n    sensor = (pl.scan_csv(f'{INPUT_DIR}/sensor_geometry.csv')\n                .with_columns([\n                    pl.col('sensor_id').cast(pl.Int16),\n                    (pl.col('sensor_id') // 60).alias('string_id'),\n                    (pl.col('sensor_id') % 60).alias('depth_id'),                \n                ])\n             )\n\n    print(sensor)\n\n#     if MODE == 'train':\n#         batch_id_start = batch_id_start\n#         batch_id_end = batch_id_start + TRAIN_N_BATCHES\n#     else:\n#         batch_id_start = meta.select(pl.col('batch_id')).collect()[0, 0]\n#         batch_id_end = meta.select(pl.col('batch_id')).collect()[-1, 0] + 1\n    sub = []\n#     for batch_id in range(batch_id_start,batch_id_end):\n    print(batch_id)\n    t = time.time()\n    max_events=TRAIN_MAX_EVENTS\n    print(f'loading batch_{batch_id}...')\n    batch = pl.scan_parquet(f'{INPUT_DIR}/{MODE}/batch_{batch_id}.parquet')\n\n    ## Limit to max_events\n    if max_events is not None:\n        batch = batch.collect()\n        last_event_id = batch.select(pl.col('event_id')).unique()[TRAIN_MAX_EVENTS-1, 0]\n        batch = batch.lazy().filter(pl.col('event_id') <= last_event_id)\n\n    ## Merge in sensor x,y,z data\n    batch = batch.join(sensor, on='sensor_id', how='left').collect()\n\n    ## The logic for auxiliary = False is very basic, we can improve on it. Initial discussion here: \n    if FIND_BEST_POINTS:\n        batch = find_best_points_polars(batch)\n\n    ## Use data point weights instead of filtering. Still set most points to 0 weight, unless they are the only data points available.\n    ## MIN_PRIMARY_DATAPOINTS and AUX_FALSE_WEIGHT are tuned values, based on optimizing the score on batches 101-110.\n    batch = batch.lazy().with_columns((pl.col('best').count().over('event_id') - pl.col('best').sum().over('event_id')).alias('best_count'))\n    batch = batch.lazy().with_columns((pl.col('auxiliary').count().over('event_id') - pl.col('auxiliary').sum().over('event_id')).alias('non_aux_count'))\n    batch = batch.with_columns(( pl.when( pl.col('best') | ((pl.col('best_count') < MIN_PRIMARY_DATAPOINTS) & (pl.col('non_aux_count') < MIN_PRIMARY_DATAPOINTS)) )\n                                            .then(1.0)\n                                            .otherwise(pl.when(pl.col('auxiliary').is_not())\n                                                         .then(AUX_FALSE_WEIGHT)\n                                                         .otherwise(0.0)\n                                                    ) \n                                      ).alias('trust'))\n    batch = batch.lazy().with_columns((pl.col('charge') * pl.col('trust')).alias('charge'))\n    batch = batch.drop('auxiliary').with_columns(pl.col(\"best\").cast(pl.Int16, strict=False).alias(\"best\"))\n    batch = time_weighted_centering(batch)\n    return batch","metadata":{"execution":{"iopub.status.busy":"2023-02-14T15:23:41.138609Z","iopub.execute_input":"2023-02-14T15:23:41.139102Z","iopub.status.idle":"2023-02-14T15:23:41.155288Z","shell.execute_reply.started":"2023-02-14T15:23:41.139069Z","shell.execute_reply":"2023-02-14T15:23:41.153992Z"},"trusted":true},"execution_count":null,"outputs":[]},{"cell_type":"code","source":"# if MODE == 'train':\n#     !rm -r train.csv\n#     integration = batch_integration(1)\n#     print('writing to csv...')\n#     integration.write_csv('train.csv')\n#     del integration\n#     print(gc.collect())\n#     detector_readings = pl.scan_csv(f'{MODE}.csv')","metadata":{"execution":{"iopub.status.busy":"2023-02-14T15:23:41.157049Z","iopub.execute_input":"2023-02-14T15:23:41.157811Z","iopub.status.idle":"2023-02-14T15:23:41.171312Z","shell.execute_reply.started":"2023-02-14T15:23:41.157760Z","shell.execute_reply":"2023-02-14T15:23:41.170197Z"},"trusted":true},"execution_count":null,"outputs":[]},{"cell_type":"code","source":"# integration","metadata":{"execution":{"iopub.status.busy":"2023-02-14T15:23:41.173006Z","iopub.execute_input":"2023-02-14T15:23:41.173425Z","iopub.status.idle":"2023-02-14T15:23:41.183167Z","shell.execute_reply.started":"2023-02-14T15:23:41.173361Z","shell.execute_reply":"2023-02-14T15:23:41.182218Z"},"trusted":true},"execution_count":null,"outputs":[]},{"cell_type":"code","source":"def add_to_table(database_path: str,\n                      df: pd.DataFrame,\n                      table_name:  str,\n                      is_primary_key: bool,\n                      ) -> None:\n    \"\"\"Writes meta data to sqlite table. \n\n    Args:\n        database_path (str): the path to the database file.\n        df (pd.DataFrame): the dataframe that is being written to table.\n        table_name (str, optional): The name of the meta table. Defaults to 'meta_table'.\n        is_primary_key(bool): Must be True if each row of df corresponds to a unique event_id. Defaults to False.\n    \"\"\"\n    try:\n        create_table(   columns= df.columns,\n                        database_path = database_path, \n                        table_name = table_name,\n                        integer_primary_key= is_primary_key,\n                        index_column = 'event_id')\n    except sqlite3.OperationalError as e:\n        if 'already exists' in str(e):\n            pass\n        else:\n            raise e\n    engine = sqlalchemy.create_engine(\"sqlite:///\" + database_path)\n#     print(type(df))\n    df = df.to_pandas()\n    \n    df.to_sql(table_name, con=engine, index=False, if_exists=\"append\", chunksize = 200000)\n    engine.dispose()\n    return\n\ndef convert_to_sqlite(meta_data_path: str,\n                      database_path: str,\n                      input_data_folder: str,\n                      batch_size: int = 200000,\n                      batch_ids: Optional[List[int]] = None,) -> None:\n    \"\"\"Converts a selection of the Competition's parquet files to a single sqlite database.\n\n    Args:\n        meta_data_path (str): Path to the meta data file.\n        batch_size (int): the number of rows extracted from meta data file at a time. Keep low for memory efficiency.\n        database_path (str): path to database. E.g. '/my_folder/data/my_new_database.db'\n        input_data_folder (str): folder containing the parquet input files.\n        batch_ids (List[int]): The batch_ids you want converted. Defaults to None (all batches will be converted)\n    \"\"\"\n    if batch_ids is None:\n        batch_ids = np.arange(1,661,1).to_list()\n    else:\n        assert isinstance(batch_ids,list), \"Variable 'batch_ids' must be list.\"\n    if not database_path.endswith('.db'):\n        database_path = database_path+'.db'\n    meta_data_iter = pq.ParquetFile(meta_data_path).iter_batches(batch_size = batch_size)\n    batch_id = 1\n    converted_batches = []\n    progress_bar = tqdm(total = len(batch_ids))\n    for meta_data_batch in meta_data_iter:\n        if batch_id in batch_ids:\n#             print(\"adding meta_data...\")\n            meta_data_batch = pl.from_pandas(meta_data_batch.to_pandas())\n            pulses = load_input(meta_batch=meta_data_batch, input_data_folder= input_data_folder)\n            whole = meta_data_batch.join(pulses, on='event_id',how=\"inner\")\n            whole = pl.from_pandas(whole.to_pandas().dropna(axis=0,how='any'))\n            meta_data_batch = whole.select(meta_data_batch.columns)\n            pulses = whole.select(pulses.columns)\n            del whole\n            gc.collect()\n#             print(meta_data_batch)\n            add_to_table(database_path = database_path,\n                        df = meta_data_batch,\n                        table_name='meta_table',\n                        is_primary_key= True)\n            \n#             print(pulses)\n            del meta_data_batch # memory\n#             print(\"meta_data added\")\n            gc.collect()\n#             print(\"adding pulse_data...\")\n            add_to_table(database_path = database_path,\n                        df = pulses,\n                        table_name='pulse_table',\n                        is_primary_key= False)\n            del pulses # memory\n#             print(\"pulse_data added\")\n            gc.collect()\n            progress_bar.update(1)\n            converted_batches.append(batch_id)\n        batch_id +=1\n        if len(batch_ids) == len(converted_batches):\n            break\n    progress_bar.close()\n    del meta_data_iter # memory\n    gc.collect()\n    print(f'Conversion Complete!. Database available at\\n {database_path}')","metadata":{"execution":{"iopub.status.busy":"2023-02-14T17:43:39.175913Z","iopub.execute_input":"2023-02-14T17:43:39.176397Z","iopub.status.idle":"2023-02-14T17:43:39.206167Z","shell.execute_reply.started":"2023-02-14T17:43:39.176344Z","shell.execute_reply":"2023-02-14T17:43:39.205031Z"},"trusted":true},"execution_count":null,"outputs":[]},{"cell_type":"code","source":"input_data_folder = '/kaggle/input/icecube-neutrinos-in-deep-ice/train'\ngeometry_table = pd.read_csv('/kaggle/input/icecube-neutrinos-in-deep-ice/sensor_geometry.csv')\nmeta_data_path = '/kaggle/input/icecube-neutrinos-in-deep-ice/train_meta.parquet'","metadata":{"execution":{"iopub.status.busy":"2023-02-14T15:23:41.204248Z","iopub.execute_input":"2023-02-14T15:23:41.204650Z","iopub.status.idle":"2023-02-14T15:23:41.235726Z","shell.execute_reply.started":"2023-02-14T15:23:41.204614Z","shell.execute_reply":"2023-02-14T15:23:41.234725Z"},"trusted":true},"execution_count":null,"outputs":[]},{"cell_type":"code","source":"!rm -r batch_529to660.db","metadata":{"execution":{"iopub.status.busy":"2023-02-14T17:43:51.722282Z","iopub.execute_input":"2023-02-14T17:43:51.722705Z","iopub.status.idle":"2023-02-14T17:43:53.044998Z","shell.execute_reply.started":"2023-02-14T17:43:51.722670Z","shell.execute_reply":"2023-02-14T17:43:53.043262Z"},"trusted":true},"execution_count":null,"outputs":[]},{"cell_type":"code","source":"# !rm -r train.csv\n# print('deleted csv...')","metadata":{"execution":{"iopub.status.busy":"2023-02-14T15:23:42.336860Z","iopub.execute_input":"2023-02-14T15:23:42.337265Z","iopub.status.idle":"2023-02-14T15:23:42.342418Z","shell.execute_reply.started":"2023-02-14T15:23:42.337224Z","shell.execute_reply":"2023-02-14T15:23:42.341183Z"},"trusted":true},"execution_count":null,"outputs":[]},{"cell_type":"code","source":"batch_ids = []\nfor i in range(529,661):\n    batch_ids.append(i)\nbatch_ids.remove(602)\nbatch_ids","metadata":{"execution":{"iopub.status.busy":"2023-02-14T15:23:42.343915Z","iopub.execute_input":"2023-02-14T15:23:42.344275Z","iopub.status.idle":"2023-02-14T15:23:42.358098Z","shell.execute_reply.started":"2023-02-14T15:23:42.344243Z","shell.execute_reply":"2023-02-14T15:23:42.356768Z"},"trusted":true},"execution_count":null,"outputs":[]},{"cell_type":"code","source":"database_path = '/kaggle/working/batch_529to660'\nconvert_to_sqlite(meta_data_path,\n                  database_path=database_path,\n                  input_data_folder=input_data_folder,\n                  batch_ids = batch_ids)","metadata":{"scrolled":true,"execution":{"iopub.status.busy":"2023-02-14T15:23:42.359841Z","iopub.execute_input":"2023-02-14T15:23:42.360199Z","iopub.status.idle":"2023-02-14T15:23:42.369410Z","shell.execute_reply.started":"2023-02-14T15:23:42.360168Z","shell.execute_reply":"2023-02-14T15:23:42.368488Z"},"trusted":true},"execution_count":null,"outputs":[]},{"cell_type":"code","source":"# database_path = '/kaggle/working/batch_529to660'\n\n# convert_to_sqlite(meta_data_path,\n#               database_path=database_path,\n#               input_data_folder=input_data_folder,\n# #               batch_size = 153924,\n#               batch_ids = [602])","metadata":{"scrolled":true,"execution":{"iopub.status.busy":"2023-02-14T17:43:54.314422Z","iopub.execute_input":"2023-02-14T17:43:54.315262Z","iopub.status.idle":"2023-02-14T17:45:35.751623Z","shell.execute_reply.started":"2023-02-14T17:43:54.315218Z","shell.execute_reply":"2023-02-14T17:45:35.750461Z"},"trusted":true},"execution_count":null,"outputs":[]},{"cell_type":"code","source":"!rm -r log\n!rm -r software","metadata":{"execution":{"iopub.status.busy":"2023-02-14T15:25:40.054864Z","iopub.status.idle":"2023-02-14T15:25:40.055324Z","shell.execute_reply.started":"2023-02-14T15:25:40.055111Z","shell.execute_reply":"2023-02-14T15:25:40.055132Z"},"trusted":true},"execution_count":null,"outputs":[]}]}