import asyncio from concurrent.futures import ThreadPoolExecutor import duckdb import pandas as pd import os def find_indicator_column(table: str, indicator_columns_per_table: dict[str,str]) -> str: """Retrieves the name of the indicator column within a table. This function maps table names to their corresponding indicator columns using the predefined mapping in INDICATOR_COLUMNS_PER_TABLE. Args: table (str): Name of the table in the database Returns: str: Name of the indicator column for the specified table Raises: KeyError: If the table name is not found in the mapping """ print(f"---- Find indicator column in table {table} ----") return indicator_columns_per_table[table] async def execute_sql_query(sql_query: str) -> pd.DataFrame: """Executes a SQL query on the DRIAS database and returns the results. This function connects to the DuckDB database containing DRIAS climate data and executes the provided SQL query. It handles the database connection and returns the results as a pandas DataFrame. Args: sql_query (str): The SQL query to execute Returns: pd.DataFrame: A DataFrame containing the query results Raises: duckdb.Error: If there is an error executing the SQL query """ def _execute_query(): # Execute the query con = duckdb.connect() HF_TOKEN = os.getenv("HF_TOKEN") con.execute(f"""CREATE SECRET hf_token ( TYPE huggingface, TOKEN '{HF_TOKEN}' );""") results = con.execute(sql_query).fetchdf() # return fetched data return results # Run the query in a thread pool to avoid blocking loop = asyncio.get_event_loop() with ThreadPoolExecutor() as executor: return await loop.run_in_executor(executor, _execute_query)