from opensearchpy import OpenSearch, RequestsHttpConnection, Search, helpers from requests_aws4auth import AWS4Auth from snowflake_executor import SFExecutor import boto3 import config def get_unresolved_es_ids_from_sf(): with SFExecutor(config.SF_PARAMS) as snowflake_executor: es_ids_tuples = snowflake_executor.get_unresolved_es_ids_from_sf( config.ENVIRONMENT) es_ids = [es_id[0] for es_id in es_ids_tuples] return es_ids def get_opensearch_client(): """Get OpenSearch client.""" region = 'us-east-1' service = 'es' credentials = boto3.Session().get_credentials() awsauth = AWS4Auth( credentials.access_key, credentials.secret_key, region, service, session_token=credentials.token ) return OpenSearch( [config.OPENSEARCH_HOST], use_ssl=config.OPENSEARCH_USE_SSL, port=config.OPENSEARCH_PORT, timeout=config.OPENSEARCH_TIMEOUT, max_retries=config.OPENSEARCH_MAX_RETRIES, retry_on_timeout=True, http_auth=awsauth, connection_class=RequestsHttpConnection) def get_es_ids_from_opensearch(): opensearch_client = get_opensearch_client() search = Search(using=opensearch_client, index=config.OPENSEARCH_INDEX) response = search.scan() es_ids = [] for h in response: es_ids.append(h.meta.id) return es_ids def set_es_ids_to_unindexed_in_snowflake(missing_from_opensearch): with SFExecutor(config.SF_PARAMS) as snowflake_executor: snowflake_executor.set_es_ids_to_unindexed(missing_from_opensearch) def delete_already_resolved_from_opensearch(es_ids): opensearch_client = get_opensearch_client() chunk_size = 500 es_ids_list = list(es_ids) for i in range(0, len(es_ids_list), chunk_size): chunk = es_ids_list[i:i + chunk_size] actions = [ { '_op_type': 'delete', '_index': config.OPENSEARCH_INDEX, '_id': es_id } for es_id in chunk ] print(f'Deleting chunk {i // chunk_size + 1}...') success, _ = helpers.bulk(opensearch_client, actions) print(f'Deleted {success} documents from opensearch.') # diff snowflake and opensearch unresolved_es_ids_in_snowflake = get_unresolved_es_ids_from_sf() all_es_ids_in_opensearch = get_es_ids_from_opensearch() missing_from_opensearch = set( unresolved_es_ids_in_snowflake) - set(all_es_ids_in_opensearch) resolved_in_sf_but_still_present_in_opensearch = set( all_es_ids_in_opensearch) - set(unresolved_es_ids_in_snowflake) print('Documents missing from opensearch:') for es_id in missing_from_opensearch: print(es_id) print('Setting these es_ids to NULL and es_indexed to FALSE in snowflake...') # update snowflake if missing_from_opensearch: # Update statement breaks for more than 32768 ids chunk_size = 32768 for i in range(0, len(missing_from_opensearch), chunk_size): chunk = list(missing_from_opensearch)[i:i + chunk_size] set_es_ids_to_unindexed_in_snowflake(chunk) print('Done.') print('Documents that have been marked as resolved in snowflake, but are still present in opensearch:') for es_id in resolved_in_sf_but_still_present_in_opensearch: print(es_id) print('Deleting these documents from opensearch...') # update opensearch if resolved_in_sf_but_still_present_in_opensearch: delete_already_resolved_from_opensearch( resolved_in_sf_but_still_present_in_opensearch) print('Done.') print('Set {} es_ids in snowflake to unindexed ' '(these will be picked up on the next run of' 'swf-yt-conflict-elasticsearch).'.format(len(missing_from_opensearch))) print('Deleted {} es_ids from opensearch.'.format( len(resolved_in_sf_but_still_present_in_opensearch)))