from dotenv import load_dotenv from pathlib import Path from snowflake import connector import uuid import pymysql load_dotenv() from config import RDS_DB_CONFIG, SNOWFLAKE_DB_CONFIG # noqa: E402 query_template = Path("query.sql").read_text() snowflake_connection = connector.connect(**SNOWFLAKE_DB_CONFIG) rds_connection = pymysql.connect( **RDS_DB_CONFIG, cursorclass=pymysql.cursors.DictCursor ) VENDOR_IDS = [ 36299, 62266, 45681, ] STATEMENT_PERIOD_ID = 292 REPORT_RUN_UUID = str(uuid.uuid4()) try: print(f"Creating report run with UUID {REPORT_RUN_UUID}") print("Fetching collaborators") with rds_connection.cursor() as cursor: cursor.execute( "select * from collaborator where vendor_id in (%s)" % ",".join(["%s"] * len(VENDOR_IDS)), tuple(VENDOR_IDS), ) collaborators = cursor.fetchall() print(f"Found {len(collaborators)} collaborators") print("Creating report entries") with rds_connection.cursor() as cursor: for collaborator in collaborators: cursor.execute( """ insert into report ( report_run_uuid, report_run_name, collaborator_id, vendor_id, period_ids, period_name, filename, source ) values (%s, %s, %s, %s, %s, %s, %s, %s) """, ( REPORT_RUN_UUID, f"Bulk test {REPORT_RUN_UUID}", collaborator["id"], collaborator["vendor_id"], STATEMENT_PERIOD_ID, f"Abacus period {STATEMENT_PERIOD_ID}", f"Bulk test {REPORT_RUN_UUID} - {collaborator['id']}", "ABACUS", ), ) rds_connection.commit() # for report in reports: print("Triggering query...", end="") query = query_template.format( report_run_uuid=REPORT_RUN_UUID, collaborator_ids=", ".join( [str(collaborator["id"]) for collaborator in collaborators] ), collaborator_ids_bracketed=", ".join( [f"({collaborator['id']})" for collaborator in collaborators] ), statement_period_id=STATEMENT_PERIOD_ID, ) with snowflake_connection.cursor(connector.DictCursor) as cursor: cursor.execute_async(query) query_id = cursor.sfqid print( " done.\nhttps://app.snowflake.com/" f"sme/orchard/#/compute/history/queries/{query_id}/detail" ) print("Updating RDS") with rds_connection.cursor() as cursor: cursor.execute( "update report set snowflake_query_id = %s where report_run_uuid = %s", (query_id, REPORT_RUN_UUID), ) rds_connection.commit() print("Done") finally: snowflake_connection.close() rds_connection.close()