"""elasticsearch_export generators.""" import gzip import json from yt_conflict_elasticsearch.flows.elasticsearch_export import config from yt_conflict_elasticsearch.flows.elasticsearch_export.logger import \ app_logger def conflict_status(context): """Generate conflict statuses.""" for conflict_status in config.CONFLICT_STATUSES: yield {'conflict_status': conflict_status} def conflicts_to_remove_os(es_ids, index_name): """Generate a list of conflicts to remove. Args: conflicts (list): a list of es_ids ids index_name (str): ES index name Yields: dict: conflict to remove """ for (es_id,) in es_ids: yield { '_index': index_name, '_op_type': 'delete', '_id': es_id } def get_unindexed_conflict_ids_by_line(gzipped_file_path): """Return an iterator to all unindexed new conflicts.""" with gzip.open(gzipped_file_path, 'rt') as f: for ln in f: conflict_data = json.loads(ln) if not conflict_data['es_id']: yield [t['conflict_id'] for t in conflict_data['territories']] def get_es_ids_to_update_territories_os(gzipped_file_path): """Return an iterator to all partially modified conflicts.""" with gzip.open(gzipped_file_path, 'rt') as f: for ln in f: conflict_data = json.loads(ln) if conflict_data['es_id']: yield { '_index': config.OS_CONFLICTS_INDEX_NAME, '_op_type': 'update', '_id': conflict_data['es_id'], 'doc': { 'territories': conflict_data['territories'] } } def get_conflict_data_to_populate_os(gzipped_file_path, index_name): """Return an iterator to conflicts that need to be inserted.""" with gzip.open(gzipped_file_path, 'rt') as f: for ln in f: conflict_data = json.loads(ln) if not conflict_data['es_id']: # Remove an empty string del conflict_data['es_id'] yield { '_index': index_name, '_source': conflict_data, } def get_batch_to_insert_to_sf_temp( conflict_ids_iterator, es_insertion_iterator): """Yield a batch of records to insert in SF temp table.""" batch_conflict_ids_with_es_ids = [] try: while True: is_success, res = next(es_insertion_iterator) conflict_ids = next(conflict_ids_iterator) # if the insertion was successful, update the list of es_ids if is_success: for conflict_id in conflict_ids: batch_conflict_ids_with_es_ids.append( {'conflict_id': conflict_id, 'es_id': res['index']['_id']}) if len(batch_conflict_ids_with_es_ids) == \ config.SF_INSERT_BATCH_SIZE: yield batch_conflict_ids_with_es_ids batch_conflict_ids_with_es_ids = [] else: app_logger.info( 'Couldn\'t insert documents ' 'for the following conflict ids: {}'.format(conflict_ids)) except StopIteration: # retrieve the remainder from our batch size if len(batch_conflict_ids_with_es_ids): yield batch_conflict_ids_with_es_ids def get_batch_to_insert_to_sf_temp_os( conflict_ids_iterator, es_insertion_iterator): """Yield a batch of records to insert in SF temp table.""" batch_conflict_ids_with_es_ids = [] try: while True: is_success, res = next(es_insertion_iterator) conflict_ids = next(conflict_ids_iterator) # if the insertion was successful, update the list of es_ids if is_success: for conflict_id in conflict_ids: batch_conflict_ids_with_es_ids.append( {'conflict_id': conflict_id, 'es_id': res['index']['_id']}) if len(batch_conflict_ids_with_es_ids) == \ config.SF_INSERT_BATCH_SIZE: yield batch_conflict_ids_with_es_ids batch_conflict_ids_with_es_ids = [] else: app_logger.info( 'Couldn\'t insert documents ' 'for the following conflict ids: {}'.format(conflict_ids)) except StopIteration: if len(batch_conflict_ids_with_es_ids): yield batch_conflict_ids_with_es_ids def insert_into_es_os(insert_operator): """Insert documents into OS and log errors.""" try: while True: try: is_success, res = next(insert_operator) if not is_success: app_logger.error( 'Couldn\'t insert document: {}'.format(res)) except StopIteration: raise StopIteration except StopIteration: pass except Exception as e: app_logger.error('Error inserting document: {}'.format(e))