from multiprocessing import Process import os import time import config from connectors import mysql logger = config.logger def run_query(sql: str) -> None: connection = mysql.get_mysql_connection() with connection.cursor() as cursor: now = time.time() cursor.execute(sql) post_query = time.time() logger.info(f"Query time for {sql} is {post_query - now}\n") result = cursor.fetchall() if result: result_log_message = f"Tracks missing ISRCs: {result}" else: result_log_message = "All tracks have assigned ISRCs" logger.info(result_log_message) if connection.show_warnings(): logger.warning(f"Warnings: {connection.show_warnings()}") connection.close() def main() -> None: database = "accountingflat" table = "TEMP_dig_sales_statements" partitions = mysql.get_table_partition_count(database, table) batch_size = config.MYSQL_CONCURRENT_LOADS for batch_start in range(0, len(partitions), batch_size): logger.info(f"Starting from partition {batch_start}") partition_batch = partitions[batch_start : batch_start + batch_size] processes = [] for partition in partition_batch: sql_file = ( os.path.realpath(os.path.dirname(__name__)) + "/accounting_run_utils/queries/validate_temp_dig_sales_statements_null_isrc.sql" ) sql = mysql.prepare_partitioned_query(sql_file, partition) processes.append(Process(target=run_query, args=([sql]))) for process in processes: process.start() for process in processes: process.join() if __name__ == "__main__": main()