{"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":"import pyarrow.parquet as pq\nimport sqlite3\nimport pandas as pd\nimport sqlalchemy\nfrom tqdm import tqdm\nimport os\nimport os.path\nfrom typing import Any, Dict, List, Optional\nimport numpy as np\n\n\n# from graphnet.data.sqlite.sqlite_utilities import create_table\n\ndef run_sql_code(database_path: str, code: str) -> None:\n    \"\"\"Execute SQLite code.\n    Args:\n        database_path: Path to databases\n        code: SQLite code\n    \"\"\"\n    conn = sqlite3.connect(database_path)\n    c = conn.cursor()\n    c.executescript(code)\n    c.close()\n    \n\ndef attach_index(\n    database_path: str, table_name: str, index_column: str = \"event_no\"\n) -> None:\n    \"\"\"Attach the table (i.e., event) index.\n    Important for query times!\n    \"\"\"\n    code = (\n        \"PRAGMA foreign_keys=off;\\n\"\n        \"BEGIN TRANSACTION;\\n\"\n        f\"CREATE INDEX {index_column}_{table_name} \"\n        f\"ON {table_name} ({index_column});\\n\"\n        \"COMMIT TRANSACTION;\\n\"\n        \"PRAGMA foreign_keys=on;\"\n    )\n    run_sql_code(database_path, code)\n\n\n\ndef create_table(\n    columns: List[str],\n    table_name: str,\n    database_path: str,\n    *,\n    index_column: str = \"event_no\",\n    default_type: str = \"NOT NULL\",\n    integer_primary_key: bool = True,\n    ) -> None:\n    \"\"\"Create a table.\n    Args:\n        columns: Column names to be created in table.\n        table_name: Name of the table.\n        database_path: Path to the database.\n        index_column: Name of the index column.\n        default_type: The type used for all non-index columns.\n        integer_primary_key: Whether or not to create the `index_column` with\n            the `INTEGER PRIMARY KEY` type. Such a column is required to have\n            unique, integer values for each row. This is appropriate when the\n            table has one row per event, e.g., event-level MC truth. It is not\n            appropriate for pulse map series, particle-level MC truth, and\n            other such data that is expected to have more that one row per\n            event (i.e., with the same index).\n    \"\"\"\n    # Prepare column names and types\n    query_columns = []\n    for column in columns:\n        type_ = default_type\n        if column == index_column:\n            if integer_primary_key:\n                type_ = \"INTEGER PRIMARY KEY NOT NULL\"\n            else:\n                type_ = \"NOT NULL\"\n\n        query_columns.append(f\"{column} {type_}\")\n    query_columns_string = \", \".join(query_columns)\n\n    # Run SQL code\n    code = (\n        \"PRAGMA foreign_keys=off;\\n\"\n        f\"CREATE TABLE {table_name} ({query_columns_string});\\n\"\n        \"PRAGMA foreign_keys=on;\"\n    )\n    run_sql_code(\n        database_path,\n        code,\n    )\n\n    # Attaching index to all non-truth-like tables (e.g., pulse maps).\n    if not integer_primary_key:\n        attach_index(database_path, table_name, index_column=index_column)\n\n\n\n\ndef load_input(meta_batch: pd.DataFrame, input_data_folder: str) -> pd.DataFrame:\n    \"\"\"\n    Will load the corresponding detector readings associated with the meta data batch.\n    \"\"\"\n    batch_id = pd.unique(meta_batch['batch_id'])\n\n    assert len(batch_id) == 1, \"contains multiple batch_ids. Did you set the batch_size correctly?\"\n        \n    detector_readings = pd.read_parquet(path = f'{input_data_folder}/batch_{batch_id[0]}.parquet')\n    sensor_positions = geometry_table.loc[detector_readings['sensor_id'], ['x', 'y', 'z']]\n    sensor_positions.index = detector_readings.index\n        \n    \n    for column in sensor_positions.columns:\n        if column not in detector_readings.columns:\n            detector_readings[column] = sensor_positions[column]\n\n    detector_readings['auxiliary'] = detector_readings['auxiliary'].replace({True: 1, False: 0})\n    return detector_readings.reset_index()\n\n\ndef 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    df.to_sql(table_name, con=engine, index=False, if_exists=\"append\", chunksize = 200000)\n    engine.dispose()\n    return\n\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    \"\"\"\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            meta_data_batch  = meta_data_batch.to_pandas()\n            add_to_table(database_path = database_path,\n                        df = meta_data_batch,\n                        table_name='meta_table',\n                        is_primary_key= True)\n            pulses = load_input(meta_batch=meta_data_batch, input_data_folder= input_data_folder)\n            del meta_data_batch # memory\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            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    print(f'Conversion Complete!. Database available at\\n {database_path}')","metadata":{"execution":{"iopub.status.busy":"2023-03-16T16:37:34.886198Z","iopub.execute_input":"2023-03-16T16:37:34.886630Z","iopub.status.idle":"2023-03-16T16:37:34.916252Z","shell.execute_reply.started":"2023-03-16T16:37:34.886592Z","shell.execute_reply":"2023-03-16T16:37:34.915257Z"},"trusted":true},"execution_count":null,"outputs":[]},{"cell_type":"code","source":"!rm '/kagge/working/batch_1.db'\ninput_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'\n\n#batch_1\ndatabase_path = 'batch_52_56.db'\nconvert_to_sqlite(meta_data_path,\n                  database_path=database_path,\n                  input_data_folder=input_data_folder,\n                  batch_ids = [52, 53, 54, 55, 56])\n\n#batch_51\n# database_path = '/kagge/working/batch_53'\n# convert_to_sqlite(meta_data_path,\n#                   database_path=database_path,\n#                   input_data_folder=input_data_folder,\n#                   batch_ids = [53])\n","metadata":{"execution":{"iopub.status.busy":"2023-03-16T16:44:43.888521Z","iopub.execute_input":"2023-03-16T16:44:43.889028Z","iopub.status.idle":"2023-03-16T16:55:20.339267Z","shell.execute_reply.started":"2023-03-16T16:44:43.888987Z","shell.execute_reply":"2023-03-16T16:55:20.338159Z"},"trusted":true},"execution_count":null,"outputs":[]}]}