import os import sys parent_dir = (os.path.dirname(os.path.abspath(__file__))) utils_dir = os.path.join(parent_dir, "utils") sys.path.append(utils_dir) import pandas as pd pd.options.mode.chained_assignment = None from sqlalchemy import create_engine, text from sqlalchemy.orm import sessionmaker import os import mypasswords class DatabaseUtils: def __init__(self, db_source=None): self.db_source = db_source if db_source else os.getenv('db_source') self.main_engine = self.create_engine(self.get_db_url('main')) self.reporting_engine = self.create_engine(self.get_db_url('reporting')) def create_engine(self, url): return create_engine(url, pool_pre_ping=True, pool_recycle=3600) def get_db_url(self, db_type): if db_type == 'main': if self.db_source == 'l': return f'postgresql+psycopg2://{mypasswords.main_db_user}:{mypasswords.main_db_password}@localhost:8015/maindb' return f'postgresql+psycopg2://{mypasswords.main_db_user}:{mypasswords.main_db_password}@main-pg-db-dev.gdbdatascience.com:5432/maindb' elif db_type == 'reporting': return f'postgresql+psycopg2://{mypasswords.reporting_db_user}:{mypasswords.reporting_db_password}@reporting-db.smeanalyticsapps.com:5439/smaredshiftdb' def query_db(self, query, db_type='main', params=None): engine = self.main_engine if db_type == 'main' else self.reporting_engine with engine.connect() as conn: df = pd.read_sql(query, conn, params=params) return df def save_to_db(self, df, schema, table, mode='append', db_type='main'): # print(f"DataFrame shape: {df.shape}") # print(f"DataFrame columns: {df.columns.tolist()}") # print(f"Schema: {schema}, Table: {table}, Mode: {mode}, DB Type: {db_type}") engine = self.main_engine if db_type == 'main' else self.reporting_engine with engine.connect() as conn: print("Saving DataFrame to the database...") df.to_sql(table, conn, schema=schema, if_exists=mode, index=False, method='multi') return df def update_db(self, query, db_type='main'): engine = self.main_engine if db_type == 'main' else self.reporting_engine with engine.begin() as conn: result = conn.execute(text(query)) return result.rowcount def delete_from_db(self, query, db_type='main'): engine = self.main_engine if db_type == 'main' else self.reporting_engine Session = sessionmaker(bind=engine) session = Session() try: with session.begin(): session.execute(query) session.commit() print("Records deleted successfully.") except Exception as e: session.rollback() print(f"Error occurred: {e}") session.close() def save_to_db_chunks(self, df, schema, table, mode='append', db_type='main'): total_rows = df.shape[0] step = 10000 print(f'Total number of rows: {total_rows}') for start in range(0, total_rows, step): chunk = df.iloc[start:start + step] self.save_to_db(chunk, schema, table, mode if start == 0 else 'append', db_type) print(f'Saved rows starting at {start}') def timedifferencer(self, datetime1, datetime2, is_overall=0): time_difference = datetime2 - datetime1 days = time_difference.days hours, remainder = divmod(time_difference.seconds, 3600) minutes, seconds = divmod(remainder, 60) if is_overall == 0: print(f"Step Time: {minutes} minutes, {seconds} seconds") else: print(f"Overall Time: {minutes} minutes, {seconds} seconds")