import os from datetime import datetime import calendar from itertools import chain from dotenv import load_dotenv load_dotenv() from utils.dates import Dates SECRETS_PATH = os.environ.get('SECRETS_PATH', 'elasticsearch_analytics') ES_INDEX = os.environ.get('ES_INDEX', '') ES_INDEX_ALIAS = os.environ.get('ES_INDEX_ALIAS', '') ES_DOMAIN = os.environ.get('ES_DOMAIN', '') AWS_ACCESS_KEY_ID = os.environ.get('AWS_ACCESS_KEY_ID', '') AWS_SECRET_ACCESS_KEY = os.environ.get('AWS_SECRET_ACCESS_KEY', '') AWS_REGION = os.environ.get('AWS_REGION', 'us-east-1') SNOWFLAKE_PRIVATE_KEY_PATH = os.environ.get('SNOWFLAKE_PRIVATE_KEY_PATH', '') SNOWFLAKE_KEY_PASSPHRASE = os.environ.get('SNOWFLAKE_KEY_PASSPHRASE', None) ENV = os.environ.get('Environment', 'dev') def merge_configs(c1, c2): """Merge two flat configs. Values from c1 get overridden by values from c2 if the keys collide. Args: c1 (dict): First config. c2 (dict): Second config. Returns: dict: Result dict containing merged result. """ c1 = c1 or {} c2 = c2 or {} return {k: v for k, v in chain(c1.items(), c2.items()) if v} # Default Snowflake connection parameters excluding credentials SF_PARAMS = { 'role': os.environ.get('SNOWFLAKE_ROLE'), 'warehouse': os.environ.get('SNOWFLAKE_WAREHOUSE'), 'db': os.environ.get('SNOWFLAKE_DATABASE'), 'schema': os.environ.get('SNOWFLAKE_SCHEMA') } # Snowflake connection credentials SF_CREDENTIALS = { 'user': os.environ.get('SNOWFLAKE_USER'), 'account': os.environ.get('SNOWFLAKE_ACCOUNT'), 'private_key': None } SF_CONFIG = merge_configs(SF_PARAMS, SF_CREDENTIALS) START_DATE = os.environ.get('START_DATE', None) END_DATE = os.environ.get('END_DATE', None) if not START_DATE and not END_DATE: month_range = Dates().get_month_range_str() START_DATE = month_range['start_date'] END_DATE = month_range['end_date'] INGEST_CONFIG = { 'start_date': START_DATE, 'end_date': END_DATE, 'chunk_limit': int(os.environ.get('CHUNK_LIMIT', 2000000)), 'offset': int(os.environ.get('OFFSET', 0)), 'max_data_limit': int(os.environ.get('MAX_DATA_LIMIT', 0)), 'no_of_processes': int(os.environ.get('NO_OF_PROCESSES', 18)), } ES_MAXSIZE = int(os.environ.get('ES_MAXSIZE', 12)) ES_TIMEOUT = int(os.environ.get('ES_TIMEOUT', 60)) ES_RETRY_ON_TIMEOUT = os.environ.get('ES_RETRY_ON_TIMEOUT', 1) == 1 ES_MAX_RETRIES = int(os.environ.get('ES_MAX_RETRIES', 1)) ES_BULK_CONFIG = { 'thread_count': int(os.environ.get('THREAD_COUNT', 5)), 'queue_size': int(os.environ.get('QUEUE_SIZE', 10)), 'max_chunk_bytes': int(os.environ.get('MAX_CHUNK_BYTES', 20971520)), 'chunk_size': int(os.environ.get('CHUNK_SIZE', 1000)), } ES_IDX_INGEST_SETTINGS = { 'number_of_shards': int(os.environ.get('INGEST_NUMBER_OF_SHARDS', 3)), 'number_of_replicas': int(os.environ.get('INGEST_NUMBER_OF_REPLICAS', 0)), 'refresh_interval': os.environ.get('INGEST_REFRESH_INTERVAL', '-1'), 'codec': os.environ.get('INGEST_CODEC', 'best_compression'), } ES_IDX_SEARCH_SETTINGS = { 'number_of_replicas': int(os.environ.get('SEARCH_NUMBER_OF_REPLICAS', 1)), 'refresh_interval': os.environ.get('SEARCH_REFRESH_INTERVAL', None), }