import sys import time import calendar from datetime import datetime from multiprocessing import Pool from utils.logger import Logger from utils.dates import Dates from connectors.elasticsearch import elasticsearch_client from connectors.snowflake import snowflake_ctx, snowflake_cursor from conf import config from methods import elasticsearch_methods as es_methods def ingest_data(start_date_str, end_date_str): chunk_limit = config.INGEST_CONFIG['chunk_limit'] offset_start = config.INGEST_CONFIG['offset'] no_of_processes = config.INGEST_CONFIG['no_of_processes'] start_date_time_obj = datetime.strptime(start_date_str, '%Y-%m-%d') start_date_time_obj = datetime.strptime(end_date_str, '%Y-%m-%d') idx_start_month = start_date_time_obj.month idx_start_year = start_date_time_obj.year index_month = f'{idx_start_year}-{idx_start_month}' index_name = f'{config.ES_INDEX}-{index_month}' es_methods.prepare_index_to_be_ingested(index_name) Logger(f'Start date {start_date_str}, End date {end_date_str}').log_info() if config.INGEST_CONFIG['max_data_limit'] > 0: data_count_to_ingest = config.INGEST_CONFIG['max_data_limit'] else: data_count_to_ingest = es_methods.fetch_data_count_from_snowflake( start_date_str, end_date_str ) offset_end = data_count_to_ingest offset_list_generator = ({ 'offset': offset, 'start_date_str': start_date_str, 'end_date_str': end_date_str } for offset in range(offset_start, offset_end, chunk_limit)) try: Logger(f'POOL_START_{offset_start}_{chunk_limit}_{offset_end}').log_info() with Pool(no_of_processes) as p: p.map(es_methods.ingest_to_es, offset_list_generator) Logger(f'POOL_END_{offset_start}_{chunk_limit}_{offset_end}').log_info() es_methods.clone_index(index_name, remove=True) es_methods.reset_index_settings() except Exception as error: Logger('pool_exception').log_exception(error) sys.exit() def ingest_by_month(start_date_str, end_date_str): start_date_time_obj = datetime.strptime(start_date_str, '%Y-%m-%d') end_date_time_obj = datetime.strptime(end_date_str, '%Y-%m-%d') if(end_date_time_obj < start_date_time_obj): Logger(f'END_DATE: {end_date_time_obj} should be >= START_DATE: {start_date_time_obj}').log_info() sys.exit() idx_start_month = start_date_time_obj.month idx_start_year = start_date_time_obj.year idx_end_month = end_date_time_obj.month idx_end_year = end_date_time_obj.year while True: month_range = Dates().get_month_range_str(year=idx_start_year, month=idx_start_month) start_date_str = month_range['start_date'] end_date_str = month_range['end_date'] Logger(f'**** START: ingest_data {start_date_str} {end_date_str} ****').log_info() ingest_data(start_date_str, end_date_str) Logger(f'**** END: ingest_data {start_date_str} {end_date_str} ****').log_info() if(idx_start_month == idx_end_month and idx_start_year == idx_end_year): break # +1 till 12 month then +1 year & reset month to 1 idx_start_month += 1 if((idx_start_month % 13) == 0): idx_start_year += 1 idx_start_month = 1 def main(): full_start = time.time() Logger('START: elasticsearch-analytics-ingest:').log_info() start_date_str = config.INGEST_CONFIG['start_date'] end_date_str = config.INGEST_CONFIG['end_date'] if not start_date_str and not end_date_str: Logger(f'Provide valid value: Start date & End date').log_info() ingest_by_month(start_date_str, end_date_str) full_seconds = time.time() - full_start Logger(f'Time for full run {str(full_seconds)} seconds.').log_info() Logger('END: elasticsearch-analytics-ingest:').log_info() snowflake_cursor.close() snowflake_ctx.close() full_seconds = time.time() - full_start Logger(f'Time for script {str(full_seconds)} seconds.').log_info() if __name__ == '__main__': main()