"""Elasticsearch export script. Related env vars: ELASTICSEARCH_HOST ELASTICSEARCH_PORT ELASTICSEARCH_USE_SSL ELASTICSEARCH_USE_AUTH ELASTICSEARCH_DEFAULT_TIMEOUT AWS_SECRET_ACCESS_KEY AWS_ACCESS_KEY_ID """ import argparse import csv import json import logging import os import tenacity from vectororder.connectors import elasticsearch from vectororder.constants import search logger = logging.getLogger(__name__) logger.setLevel(getattr(logging, os.environ.get('LOGGING_LEVEL', 'INFO'))) logger.addHandler(logging.StreamHandler()) ITEMS_PER_REQUEST = int(os.environ.get('ITEMS_PER_REQUEST') or 10000) # Note: ELASTICSEARCH_DEFAULT_TIMEOUT env var respected by the connector. RETRY_ATTEMPTS = int(os.environ.get('RETRY_ATTEMPTS') or 10) RETRY_MAX_WAIT_TIME = int(os.environ.get('RETRY_MAX_WAIT_TIME') or 120) FIELD_NAMES = [ 'delivery_started', 'delivery_ended', 'encoding_started', 'encoding_ended', 'encoding_queue_detail_id', 'encoding_queue_id', 'status', 'upc', 'order_id', 'order_type', 'priority', 'meta_update', 'encoder_id', 'created_at', 'user_id', 'encoding', 'delivery', 'display_upc', 'product_id', 'store_id', ] def main(*, file_name, **kwargs): """Export results to CSV file. Args: file_name (str): File name to save CSV into. kwargs: Keyword arguments passed to the get_es_results(). """ with open(file_name, 'w', newline='\n') as csv_file: writer = csv.DictWriter( csv_file, fieldnames=FIELD_NAMES, lineterminator='\n') writer.writeheader() for results in get_es_results(**kwargs): writer.writerows(get_dicts_from_results(results)) csv_file.flush() def get_dicts_from_results(results): """Get dictionary with required fields from Opensearch request results. Args: results (iterable): Elasticsearch results. Yields: dict: A single result record with a subset of fields. """ for result in results: source = result['_source'] for k, v in source.items(): if isinstance(v, bool): if v is True: source[k] = 'true' else: source[k] = 'false' yield {f: source[f] for f in FIELD_NAMES if f in source} def get_es_results(field_name, min_value=None, max_value=None): """Fetch Elasticsearch data in chunks. Args: field_name (str): Field name to use in 'search_after' requests, e.g. encoding_queue_detail_id. min_value (int): Minimum ID to export. max_value (int): Maximum ID to export. Yields: list: Results of the Elasticsearch request with up to ITEMS_PER_REQUEST values. """ @tenacity.retry( stop=tenacity.stop_after_attempt(RETRY_ATTEMPTS), wait=tenacity.wait_exponential(multiplier=2, max=RETRY_MAX_WAIT_TIME), before_sleep=tenacity.before_sleep_log(logger, logging.WARNING), reraise=True, ) def _retry_es_request(url, payload_json): return elasticsearch.client.transport.perform_request( method='POST', url=url, body=payload_json) index_name = search.OS_VO_DETAIL_INDEX_NAME url = f'/{index_name}/_search' search_after = None iterations = 0 items_fetched = 0 while True: payload = { 'sort': field_name, 'query': {'range': {field_name: {'lte': max_value}}}, 'size': ITEMS_PER_REQUEST, } if search_after is not None: payload['search_after'] = search_after if min_value and not search_after: payload['query']['range'][field_name].update({'gte': min_value}) logger.debug( f'Fetching {ITEMS_PER_REQUEST} items with {max_value:d} >= ' f'{field_name:s} >= {min_value}, search_after = {search_after}') result = _retry_es_request(url, json.dumps(payload)) hits = len(result['hits']['hits']) iterations += 1 items_fetched += hits logger.debug( 'iteration: %d, total items fetched: %d', iterations, items_fetched) yield result['hits']['hits'] if hits < ITEMS_PER_REQUEST: logger.info(f'No more results available.') break search_after = result['hits']['hits'][-1]['sort'] def setup_parser(): """Setup args parser.""" parser = argparse.ArgumentParser() parser.add_argument( '--field-name', required=False, default='encoding_queue_detail_id', dest='field_name', help='Field name to use in \'search_after\' requests, ' 'e.g. encoding_queue_detail_id.') parser.add_argument( '--out-filename', required=False, default='es_export.csv', dest='out_filename', help='Output filename.') parser.add_argument( '--min-value', required=False, default=0, dest='min_value', type=int, help='Minimum ID to export. Only integer is supported.') parser.add_argument( '--max-value', required=False, default=None, dest='max_value', type=int, help='Maximum ID to export. Only integer is supported.') return parser if __name__ == '__main__': parser = setup_parser() args = parser.parse_args() main( field_name=args.field_name, file_name=args.out_filename, max_value=args.max_value, min_value=args.min_value, )