from multiprocessing import Process import config from connectors import mysql logger = config.logger def main() -> None: """ Get the partition count of dig_sales_detail. Then for each partition, select into outfile to S3. """ database = "accountingflat" table = "dig_sales_detail" 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: object_key = f"{config.S3_BUCKET}/{config.S3_OBJECT_PREFIX}/temp_dig_sales_statements/tmp_dig_sales_{config.PERIOD_ID}_{partition}.txt" sql = f"SELECT dsd.*, IF(dsd.retail_price IS NULL OR dsd.retail_price = 0, 2*dsd.total, dsd.qty*dsd.retail_price) AS retail , ds.period_id , cm.customer_id , cm.name AS dms_name , v.vendor_id , v.owner , o.owner_id , ds.activity_rate , ds.original_currency_id , r.release_status INTO OUTFILE S3 's3://{object_key}' FIELDS ESCAPED BY '\\\\' TERMINATED BY '\\t' OPTIONALLY ENCLOSED BY '\"' LINES TERMINATED BY '\\n' FROM `{database}`.`{table}` PARTITION ({partition}) dsd LEFT JOIN `accountingflat`.`processed_dig_sales` pds ON pds.statement_detail_id = dsd.statement_detail_id AND pds.period_id = { config.PERIOD_ID} LEFT JOIN art_relations.dig_sales_errors dse ON dse.statement_detail_id = dsd.statement_detail_id AND dse.period_id = {config.PERIOD_ID} INNER JOIN art_relations.dig_sales ds ON dsd.statement_id = ds.statement_id INNER JOIN art_relations.customer_master cm ON cm.customer_id = ds.dms_customer_id LEFT JOIN art_relations.releases r ON dsd.upc = r.upc LEFT JOIN art_relations.artist_info a ON a.artist_id = r.artist_id LEFT JOIN art_relations.vendor v ON a.vendor_id = v.vendor_id LEFT JOIN art_relations.`owner` o ON o.owner_abbrivation = v.owner WHERE ds.paid = 'Y' AND ds.period_id = {config.PERIOD_ID} AND pds.statement_detail_id IS NULL AND dse.statement_detail_id IS NULL GROUP BY dsd.statement_detail_id;" processes.append( Process( target=mysql.select_data_into_outfile, args=(sql, object_key), ) ) logger.info(sql) for process in processes: process.start() for process in processes: process.join() if __name__ == "__main__": main()