import csv import json import logging import os import sys import time from functools import partial from multiprocessing import Pool import constants as const import requests QA_ENVIRONMENT = 'qa' ENVIRONMENT = os.environ.get('Environment', QA_ENVIRONMENT) MAX_PROCESSES = os.environ.get('MAX_PROCESSES', 20) FILENAME = os.environ.get('FILENAME', 'data/performers.csv') sys.path.append(os.path.join(os.path.dirname(__file__), '..')) log = logging.getLogger('main') log.setLevel(logging.DEBUG) fmt = logging.Formatter('%(asctime)s - %(levelname)s - %(message)s') sh = logging.StreamHandler(sys.stdout) sh.setFormatter(fmt) log.addHandler(sh) log.info('Running performer bulkupdate for %s Environment with %s workers', ENVIRONMENT, MAX_PROCESSES) DOMAIN_PREFIX = 'qa-' if ENVIRONMENT == 'prod': DOMAIN_PREFIX = 'prod-' def get_roles(): """ Get available performer roles from ows-track. Returns: list: list of dicts """ url = const.OWS_TRACK.format(DOMAIN_PREFIX) + const.ROLES_URL response = requests.get(url) if response.status_code != 200: log.error('Error: while ows-track GET roles request %s', response.json()) roles = {} for item in response.json()['items']: roles[item['performer_role']] = item['performer_role_id'] return roles def put_performer(track_id, data): """ Request to ows-track for PUT performers. ows-track sample PUT payload format { "birth_name" : John, "performer_role_id" : 12, "type" : "Primary Performer" } Args: track_id (int): unique identifier of track. data (dict): row of CSV file. """ log.info('ows-track request for track_id %s with payload %s', track_id, json.dumps(data)) url = const.OWS_TRACK.format(DOMAIN_PREFIX) \ + const.PERFORMER_URL.format(track_id) headers = { 'Content-Type': 'application/json', } response = requests.put(url, headers=headers, data=json.dumps(data)) if response.status_code != 200: log.error('Error: while ows-track PUT request track_id: %s', track_id) log.info('Updated performer for track_id: %s', track_id) def parse_performers(data, roles): """ Parse performer data from each row. Args: data (dict): row of CSV file. roles (dict): valid roles exists in database. Returns: list: list of dicts """ performer_artists = [] for i in range(0, 13): artist = {} post_fix = '' if not i else str(i) performer_type = const.PERFORMER_TYPE.format(post_fix) performer_name = 'Performer Legal/Birth Name' + post_fix performer_role = const.PERFORMER_ROLE.format(post_fix) if performer_type in data.keys() and data[performer_type] != '': role_id = get_role_id(roles, data[performer_role], data[const.TRACK_ID]) if not role_id: continue artist['type'] = const.VALID_PERFORMER_TYPE[data[performer_type]] artist['birth_name'] = data[performer_name] artist['performer_role_id'] = role_id performer_artists.append(artist) return performer_artists def get_role_id(roles, role_name, track_id): """ Get the numeric role_id against a role string. Args: roles (dict): valid roles exists in database. role_name (str): role name from CSV file. track_id (int): unique identifier of track Returns: int: unique role_id """ if role_name not in roles: log.error('Error: role does not exists. track_id: %s and role: %s', track_id, role_name) return False return roles[role_name] def _update_performers(track): """ Make request to put performers. Args: track (dict): track metadata Returns: int: unique role_id """ if not track: return False for track_id in track: return put_performer(track_id, {'performers': track.get(track_id)}) def main(): """Read performers data from CSV file.""" log.info('Reading file: : %s', FILENAME) roles = get_roles() performers = [] with open(FILENAME) as csv_file: reader = csv.DictReader(csv_file) for row in reader: track_id = row[const.TRACK_ID] if not track_id: continue performer_artists = [] performer_artists = parse_performers(row, roles) performers.append({track_id: performer_artists}) start = time.time() with Pool(processes=int(MAX_PROCESSES)) as pool: pool.map( partial( _update_performers), performers) end = time.time() log.info('Total processing time : %s', (end-start)) if __name__ == '__main__': main()