import csv import sys import os import psutil from connector import mysql import logging from sql import tracks, product as product_sql import math from datetime import datetime from util import formatter import boto3 import config 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) FILENAME = os.environ.get('FILENAME', 'yt_metadata.csv') pid = os.getpid() py = psutil.Process(pid) memory_use = py.memory_info()[0] print('memory use:', memory_use) session = boto3.Session() s3 = session.resource('s3') def main(): log_path = 'logs' does_log_path_exist = os.path.exists(log_path) if not does_log_path_exist: os.makedirs(log_path) WRITE_FILENAME = log_path + '/video_claimed_isrc_{}.csv'.\ format(datetime.now().strftime('%Y%m%d%H%M%S')) ar_db = None cursor = None ar_db = mysql.get_ar_mysql_connection() cursor = ar_db.cursor() with open(FILENAME, newline='', encoding='utf-8') as csvfile: reader = csv.DictReader(csvfile, delimiter=',', quotechar='"') with open(WRITE_FILENAME, "w", encoding='utf-8') as fp: writer = csv.DictWriter(fp, fieldnames=formatter.FILE_COLUMNS) writer.writeheader() for row in reader: if len(row['ISRC'].strip()) > 0: continue channel_id = row['CHANNEL_ID'] track_name = row['VIDEO_TITLE'].encode('ascii', errors='ignore') video_id = row['VIDEO_ID'] asset_id = row['ASSET_ID'] custom_id = row['CUSTOM_ID'] content_id_matching = row['CONTENT_ID_MATCHING'] auto_claim = row['AUTO_CLAIM'] monetization_status = row['MONETIZATION_STATUS'] youtube_channel_asset_type = row['YOUTUBE_CHANNEL_ASSET_TYPE'] length_minute = math.floor(int(row['VIDEO_DURATION_SEC']) / 60) length_seconds = int(row['VIDEO_DURATION_SEC']) % 60 cursor.execute(product_sql.SQL_GET_UPC.format(channel_id)) products = cursor.fetchall() for product in products: upc = product.get('upc') #I assumed track.p_line = product c_line p_line = product.get('c_line') release_id = product.get('release_id') if release_id: cursor.execute(tracks.SQL_CLAIM_ISRC) isrc_obj = cursor.fetchone() claimed_isrc = isrc_obj.get('isrc') cursor.execute(tracks.SQL_GET_TRACK_ID_BY_UPC.format(upc)) track = cursor.fetchone() track_id = track.get('track_id') cursor.execute(tracks.SQL_INSERT_TRACK, (track_id, upc, release_id, track_name, claimed_isrc, length_minute, length_seconds, p_line)) unique_track_id = cursor.lastrowid ar_db.commit() raw_record = {'channel_id': channel_id, 'upc': upc, 'unique_track_id': unique_track_id, 'track_name': track_name, 'length_minute': length_minute, 'length_seconds': length_seconds, 'isrc': claimed_isrc, 'artist_name': product.get('artist_name'), 'genre': product.get('genre'), 'video_id': video_id, 'asset_id': asset_id, 'custom_id': custom_id, 'content_id_matching': content_id_matching, 'auto_claim': auto_claim, 'monetization_status': monetization_status, 'youtube_channel_asset_type': youtube_channel_asset_type } log.info(raw_record) writer.writerow(raw_record) s3.Bucket(config.S3_BUCKET).upload_file(WRITE_FILENAME, WRITE_FILENAME) ar_db.close() memoryUse = py.memory_info()[0] print('memory use:', memoryUse) if __name__ == "__main__": main()