import config import queries from subprocess import Popen import sys from sqlalchemy import create_engine from sqlalchemy.orm import sessionmaker from sqlalchemy.pool import NullPool NUM_PROCESSES = int(sys.argv[1]) NUMBER_OF_TRACKS_TO_UPDATE_PER_TX = int(sys.argv[2]) SLEEP_TIME_PER_UPDATE = int(sys.argv[3]) _art_relations_engine = create_engine( config.DB_URL, poolclass=NullPool) _art_relations_session = sessionmaker(bind=_art_relations_engine) session = _art_relations_session() count_rows = session.execute(queries.total_track_count_query) track_count = 0 for row in count_rows: track_count = row['track_count'] break print('total number of tracks {}'.format(track_count)) ROWS_PER_PROCESS = int(track_count / NUM_PROCESSES) processes = [] for process_number in range(NUM_PROCESSES): LIMIT = ROWS_PER_PROCESS START_ROW = process_number * LIMIT END_ROW = ROWS_PER_PROCESS + START_ROW cmd = 'python run.py {NUMBER_OF_TRACKS_TO_UPDATE_PER_TX} {START_ROW} {END_ROW} {process_number} {SLEEP_TIME_PER_UPDATE}'.format( NUMBER_OF_TRACKS_TO_UPDATE_PER_TX=NUMBER_OF_TRACKS_TO_UPDATE_PER_TX, START_ROW=START_ROW, END_ROW=END_ROW, process_number=process_number, SLEEP_TIME_PER_UPDATE=SLEEP_TIME_PER_UPDATE) print(cmd) process = Popen(cmd, shell=True) processes.append(process) remaining_tracks = track_count % NUM_PROCESSES if remaining_tracks: cmd = 'python run.py {NUMBER_OF_TRACKS_TO_UPDATE_PER_TX} {START_ROW} {END_ROW} {process_number} {SLEEP_TIME_PER_UPDATE}'.format( NUMBER_OF_TRACKS_TO_UPDATE_PER_TX=NUMBER_OF_TRACKS_TO_UPDATE_PER_TX, START_ROW=track_count - remaining_tracks, END_ROW=track_count, process_number=len(processes), SLEEP_TIME_PER_UPDATE=SLEEP_TIME_PER_UPDATE) print(cmd) process = Popen(cmd, shell=True) processes.append(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()