krinya's picture
Enhance Sales Assistant with Pydantic Validation and Gradio Dashboard
e252f82
Raw
History Blame Contribute Delete
2.26 kB
"""
Database utilities for querying and writing to MySQL through an SSH tunnel.
"""
from dotenv import load_dotenv
load_dotenv()
import pandas as pd
from sqlalchemy import text
from .db_connections import connect_ssh_tunnel, create_db_engine, disconnect
def read_sql(query_str: str) -> pd.DataFrame:
"""
Execute a read SQL query and return the results as a pandas DataFrame.
Uses pandas built-in read_sql_query method.
Args:
query_str: SQL query string to execute
Returns:
pandas DataFrame with query results
Raises:
Exception: Re-raises any database or connection errors for proper error handling
"""
tunnel = None
engine = None
try:
tunnel = connect_ssh_tunnel()
engine = create_db_engine(tunnel)
df = pd.read_sql_query(sql=text(query_str), con=engine)
return df
except Exception as e:
print(f"Error executing read query: {e}")
# Re-raise the exception so it can be handled by the calling function
raise e
finally:
disconnect(engine, tunnel)
def write_to_table(df: pd.DataFrame, table_name: str, if_exists: str = 'append', schema: str = None) -> bool:
"""
Write a pandas DataFrame to a MySQL table using pandas built-in to_sql method.
Args:
df: pandas DataFrame to write
table_name: Name of the target table
if_exists: What to do if table exists ('append', 'replace', 'fail')
schema: Database schema name (optional)
Returns:
True if successful, False otherwise
"""
tunnel = None
engine = None
try:
tunnel = connect_ssh_tunnel()
engine = create_db_engine(tunnel)
df.to_sql(
name=table_name,
con=engine,
if_exists=if_exists,
index=False,
schema=schema,
method='multi',
chunksize=1000
)
print(f"Successfully wrote {len(df)} rows to table '{table_name}' (mode: {if_exists})")
return True
except Exception as e:
print(f"Error writing DataFrame to table: {e}")
return False
finally:
disconnect(engine, tunnel)