armanddemasson's picture
feat: updated common talk to data for talk to ipcc and drias
983a080
Raw History Blame
1.93 kB
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)