import logging import os import mysql.connector from dotenv import load_dotenv load_dotenv() logger = logging.getLogger(__name__) logging.basicConfig( level=logging.INFO, format='%(asctime)s %(levelname)s [%(name)s]: %(message)s' ) SOURCE_MYSQL_HOST = os.environ.get('SOURCE_MYSQL_HOST') SOURCE_MYSQL_USER = os.environ.get('SOURCE_MYSQL_USER') SOURCE_MYSQL_PASSWORD = os.environ.get('SOURCE_MYSQL_PASSWORD') DEST_MYSQL_HOST = os.environ.get('DEST_MYSQL_HOST') DEST_MYSQL_USER = os.environ.get('DEST_MYSQL_USER') DEST_MYSQL_PASSWORD = os.environ.get('DEST_MYSQL_PASSWORD') BATCH_ID = os.environ.get('BATCH_ID') def main(): try: logger.info(f'Starting copy_data for batch {BATCH_ID} from ' f'{SOURCE_MYSQL_HOST} to {DEST_MYSQL_HOST}') source_connection = mysql.connector.connect( host=SOURCE_MYSQL_HOST, port=3306, user=SOURCE_MYSQL_USER, password=SOURCE_MYSQL_PASSWORD) dest_connection = mysql.connector.connect( host=DEST_MYSQL_HOST, port=3306, user=DEST_MYSQL_USER, password=DEST_MYSQL_PASSWORD) source_cursor = source_connection.cursor() dest_cursor = dest_connection.cursor() source_query = f'SELECT * FROM sales_file_delivery.dig_sales_testfile_abacus WHERE batch_id = {BATCH_ID}' dest_query = f'INSERT INTO sales_file_delivery.dig_sales_testfile_abacus VALUES (%s, %s, %s, %s, %s, %s, %s, %s, %s, %s, %s, %s, %s, %s, %s, %s, %s, %s, %s, %s, %s, %s, %s, %s, %s, %s, %s, %s, %s, %s)' source_cursor.execute(source_query) while True: rows = source_cursor.fetchmany(size=10000) if len(rows) == 0: break logger.info(f'Fetched {len(rows)} rows from the source') dest_cursor.executemany(dest_query, rows) dest_connection.commit() logger.info(f'Inserted {dest_cursor.rowcount} rows to the destination') source_cursor.close() dest_cursor.close() source_connection.close() dest_connection.close() logger.info('Done copying data') except Exception as e: logger.error(f'Error in copy_data: {e}') raise e if __name__ == '__main__': main()