from subprocess import Popen import sys import config from connectors import mysql import query LOCK_TERRITORIES = config.LOCK_TERRITORIES with mysql.mr_session_scope() as session: # Truncate error table session.execute(query.truncate_error_table) row_count = session.execute( query.count_backfill_query) for count in row_count: total_rows = count[0] NUM_PROCESSES = 20 BATCH_DIVISOR = 10 ROWS_PER_PROCESS = int(total_rows / NUM_PROCESSES) MOD_ROWS_PER_PROCESS = int(total_rows % NUM_PROCESSES) BATCH_SIZE = int(ROWS_PER_PROCESS / BATCH_DIVISOR) MOD_BATCH_SIZE = int(ROWS_PER_PROCESS % BATCH_DIVISOR) processes = [] for i in range(NUM_PROCESSES): LIMIT = ROWS_PER_PROCESS OFFSET = i*LIMIT cmd = 'python populate_backfill_queue.py {0} {1} {2} {3} {4}'.format( LIMIT, OFFSET, BATCH_SIZE, MOD_BATCH_SIZE, LOCK_TERRITORIES) process = Popen(cmd, shell=True) processes.append(process) LIMIT = MOD_ROWS_PER_PROCESS OFFSET = NUM_PROCESSES*ROWS_PER_PROCESS cmd = 'python populate_backfill_queue.py {0} {1} {2} {3} {4}'.format( LIMIT, OFFSET, BATCH_SIZE, MOD_BATCH_SIZE, LOCK_TERRITORIES) last_process = Popen(cmd, shell=True) processes.append(last_process) exit_codes = [p.wait() for p in processes] # if any processes failed return non zero exit status if sum(exit_codes): sys.exit(1) sys.exit()