from multiprocessing import Process import os import config from connectors import mysql logger = config.logger 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/select_temp_dig_sales_processed.sql" ) sql = mysql.prepare_partitioned_query(sql_file, partition) processes.append( Process(target=mysql.run_query_with_no_result, args=([sql])) ) for process in processes: process.start() for process in processes: process.join() if __name__ == "__main__": main()