import os import ijson import json import smart_open from snowflake_connector.etl_connector import SnowflakeSQLExecutor sf_config = { 'role': os.environ.get('SNOWFLAKE_ROLE'), 'warehouse': os.environ.get('SNOWFLAKE_WAREHOUSE'), 'db': os.environ.get('SNOWFLAKE_DATABASE'), 'schema': os.environ.get('SNOWFLAKE_SCHEMA'), 'user': os.environ.get('SNOWFLAKE_USER'), 'password': os.environ.get('SNOWFLAKE_PASSWORD'), 'account': os.environ.get('SNOWFLAKE_ACCOUNT') } s3_dir = 's3://dev-aggregate-load-vbogatyrev/gracenote/' def iter_objects(s3_filename, obj_path): with smart_open.smart_open(s3_filename, 'rb') as fin: for artist in ijson.items(fin, obj_path): yield json.dumps(artist) def paginate(elems, page_size=3000): page = [] for e in elems: page.append(e) if len(page) >= page_size: yield page page = [] if len(page) > 0: yield page def load_table(table_name, s3_file, obj_path): with SnowflakeSQLExecutor(sf_config) as executor: cnt = 0 s3_path = s3_dir + s3_file for batch in paginate(iter_objects(s3_path, obj_path)): # if cnt > 9: # break # print(json.dumps(batch, indent=4)) print('Batch: %s' % (cnt,)) executor.executemany(""" INSERT INTO {} SELECT PARSE_JSON(Column1) AS data FROM VALUES (%s) """.format(table_name), batch) cnt += 1 def main(): # load_table( # 'gracenote_artist_variant', # 'gracenote_artist_20190222.json.gz', # 'artists.item') # load_table( # 'gracenote_album_master_variant', # 'gracenote_album_master_20190222.json.gz', # 'masters.item') # load_table( # 'gracenote_album_edition_variant', # 'gracenote_album_edition_20190222.json.gz', # 'editions.item') # load_table( # 'gracenote_recording_variant', # 'gracenote_recording_20190222.json.gz', # 'recordings.item') # load_table( # 'gracenote_song_variant', # 'gracenote_song_20190222.json.gz', # 'songs.item') # load_table( # 'gracenote_descriptor_variant', # 'gracenote_descriptor_20190222.json.gz', # 'descriptors') load_table( 'gracenote_hierarchy_variant', 'gracenote_hierarchy_20190222.json.gz', 'hierarchies.item') if __name__ == '__main__': main()