from datetime import datetime from json import dump, dumps, load, loads from os import getenv from os.path import isfile from boto3 import Session from botocore.exceptions import ProfileNotFound from elasticsearch import Elasticsearch, RequestsHttpConnection from elasticsearch.exceptions import RequestError from elasticsearch.helpers import bulk from elasticsearch_dsl import Search from mysql.connector import connect from requests_aws4auth import AWS4Auth from env import ENV today = datetime.today() def get_client(profile_name="default"): service = "es" host = ENV["elasticsearch_host"] region = "us-east-1" try: credentials = Session(profile_name=profile_name).get_credentials() access_key = credentials.access_key secret_key = credentials.secret_key except ProfileNotFound: access_key = ENV["aws_access_key"] secret_key = ENV["aws_secret_key"] awsauth = AWS4Auth(access_key, secret_key, region, service) return Elasticsearch( hosts=[{"host": host, "port": 443}], http_auth=awsauth, use_ssl=True, verify_certs=True, connection_class=RequestsHttpConnection ) def create_index(client, index_name): try: elasticsearch_client.indices.create( index=index_name, body={ "settings": { "analysis": { "analyzer": { "orchard_analyzer": { "tokenizer": "orchard_tokenizer", "filter": [ "lowercase", "asciifolding" ] } }, "tokenizer": { "orchard_tokenizer": { "type": "edge_ngram", "min_gram": 1, "max_gram": 10, "token_chars": [ "digit", "letter" ] } } } }, "mappings": { "_doc": { "properties": { "amg_id": { "type": "keyword" }, "id": { "type": "keyword" }, "itunes_id": { "type": "keyword" }, "name": { "analyzer": "orchard_analyzer", "search_analyzer": "orchard_analyzer", "type": "text" }, "rai_id": { "type": "keyword" }, "vendor_ids": { "type": "integer" } } } } } ) except RequestError as e: if e.error != 'resource_already_exists_exception': raise print('Skipping creating index because it already exists.') def get_elasticsearch_snapshot_filename(): return "elasticsearch_snapshot_{}_{}_{}.dat".format( today.year, today.month, today.day ) def get_elasticsearch_payloads_filename(): return "elasticsearch_payloads_{}_{}_{}.json".format( today.year, today.month, today.day ) def dump_index(client, index_name): snapshot_filename = get_elasticsearch_snapshot_filename() if isfile(snapshot_filename): print('Skipping dump because it already exists.') return search_cursor = Search(using=client, index=index_name) print("DUMP START", datetime.now()) print("Document Count: {}".format(client.cat.count())) with open(snapshot_filename, "w", encoding="utf8") as snapshot_file: for document in search_cursor.scan(): print(dumps(document.to_dict()), file=snapshot_file) print("DUMP FINISH", datetime.now()) def build_payloads(): payloads_filename = get_elasticsearch_payloads_filename() if isfile(payloads_filename): print('Skipping payload build because it already exists.') return payloads = {} print("BUILD PAYLOAD START", datetime.now()) elasticsearch_snapshot_filename = get_elasticsearch_snapshot_filename() # Mapping of rai_id to pk. unique_artist_snapshot = load(open("temp_unique_artist_snapshot.json", "r")) with open(elasticsearch_snapshot_filename, "r") as elasticsearch_snapshot_file: for document in elasticsearch_snapshot_file.readlines(): payload = loads(document) rai_id = str(payload["rai_id"]) pk = unique_artist_snapshot[rai_id] payload["id"] = pk payloads[pk] = payload dump(payloads, open(payloads_filename, "w"), indent=4) print("BUILD PAYLOAD FINISH", datetime.now()) def upload_payloads(client, index_name): payloads_filename = get_elasticsearch_payloads_filename() print("UPLOADING PAYLOADS START", datetime.now()) payloads = load(open(payloads_filename, "r")) actions = ( { "_index": index_name, "_type": "_doc", "_id": key, "_source": payload } for key, payload in payloads.items() ) bulk(client, actions) print("PLOADING PAYLOADS FINISH", datetime.now()) if __name__ == "__main__": index_name = "artists" elasticsearch_client = get_client() create_index(elasticsearch_client, index_name) dump_index(elasticsearch_client, index_name) build_payloads() upload_payloads(elasticsearch_client, index_name) # elasticsearch_client.indices.delete(index=index_name)