from elasticsearch import Elasticsearch from elasticsearch_dsl import Search from elasticsearch.helpers import parallel_bulk as elasticsearch_bulk from snowflake_executor import SFExecutor import config def get_es_ids_from_snowflake(): with SFExecutor(config.SF_PARAMS) as snowflake_executor: es_ids_tuples = snowflake_executor.get_es_ids() es_ids = [es_id[0] for es_id in es_ids_tuples] return es_ids def get_elasticsearch_client(): """Get Elasticsearch client.""" return Elasticsearch( [config.ELASTICSEARCH_HOST], use_ssl=config.ELASTICSEARCH_USE_SSL, port=config.ELASTICSEARCH_PORT, timeout=config.ELASTICSEARCH_TIMEOUT, max_retries=config.ELASTICSEARCH_MAX_RETRIES, retry_on_timeout=True) def get_es_ids_from_elasticsearch(): elasticsearch_client = get_elasticsearch_client() s = Search(using=elasticsearch_client, index='conflicts', doc_type='document') res = s.scan() es_ids = [h.meta.id for h in res] return es_ids sf_es_ids = get_es_ids_from_snowflake() es_ids = get_es_ids_from_elasticsearch() missing_from_elasticsearch = set(sf_es_ids) - set(es_ids) missing_from_snowflake = set(es_ids) - set(sf_es_ids) print('Documents missing from elasticsearch:') for es_id in missing_from_elasticsearch: print(es_id) print('Documents missing from snowflake:') for es_id in missing_from_snowflake: print(es_id)