import json import sys from connectors import mysql from connectors import sqs import db_constant import query LIMIT = int(sys.argv[1]) OFFSET = int(sys.argv[2]) BATCH_SIZE = int(sys.argv[3]) MOD_BATCH_SIZE = int(sys.argv[4]) IS_LOCKED = int(sys.argv[5]) errors = [] territories_standard = 'ISO_3166_1_2016' def queue_rows(rows): for row in rows: message_body = { db_constant.ID: row[db_constant.ID], db_constant.FILENAME: row[db_constant.FILENAME], db_constant.ISRC: row[db_constant.ISRC], db_constant.IS_LOCKED: IS_LOCKED, db_constant.DERIVED_TUID: row[db_constant.DERIVED_TUID], db_constant.DERIVED_OWNERSHIP: row[db_constant.DERIVED_OWNERSHIP]} # send to SQS sqs_response = sqs.mr_backfill.send_message( MessageBody=json.dumps(message_body)) # save SQS errors if sqs_response['ResponseMetadata']['HTTPStatusCode'] != 200: errors.append(message_body) with mysql.mr_session_scope() as session: ORIGINAL_OFFSET = OFFSET while OFFSET < ORIGINAL_OFFSET + LIMIT - MOD_BATCH_SIZE: rows = session.execute( query.select_backfill_query, {'limit': BATCH_SIZE, 'offset': OFFSET}) queue_rows(rows) OFFSET += BATCH_SIZE rows = session.execute( query.select_backfill_query, {'limit': MOD_BATCH_SIZE, 'offset': OFFSET}) queue_rows(rows) with mysql.mr_session_scope() as session: # write to error table if any errors if errors: for error in errors: session.execute(query.insert_error, error)