"""elasticsearch_export workflow config.""" import base64 from logging import INFO import os ENVIRONMENT = os.environ.get('Environment', 'dev') SWF_FLOW_NAME = 'elasticsearch_export' SWF_FLOW_VERSION = '1.0' SCRIPT_NAME = 'swf-yt-conflict-elasticsearch' 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'), 'user': os.environ.get('SNOWFLAKE_USER'), 'private_key': base64.b64decode(bytes(os.environ.get('SNOWFLAKE_KEY'), encoding='utf-8')), # noqa 'account': os.environ.get('SNOWFLAKE_ACCOUNT'), 'ar_db': os.environ.get('SNOWFLAKE_AR_DATABASE', 'ORCHARD_APP_REPORTING_V2'), # noqa 'ar_schema': os.environ.get('SNOWFLAKE_AR_SCHEMA', 'ART_RELATIONS_PROD_ART_RELATIONS'), # noqa 'ows_conflict_db': os.environ.get( 'SNOWFLAKE_OWS_CONFLICT_DATABASE', 'OWS_CONFLICT_MANAGER'), 'ows_conflict_schema': os.environ.get('SNOWFLAKE_SCHEMA'), 'temp_conflict_es_id_table': os.environ.get( 'TEMP_CONFLICT_ES_ID_TABLE', 'temp_conflict_to_es_id') } AWS_ACCESS_KEY_ID = os.environ.get('AWS_ACCESS_KEY_ID') AWS_SECRET_ACCESS_KEY = os.environ.get('AWS_SECRET_ACCESS_KEY') OPENSEARCH_HOST = os.environ.get('OPENSEARCH_HOST') OPENSEARCH_PORT = int(os.environ.get('OPENSEARCH_PORT', 443)) OPENSEARCH_USE_SSL = bool( int(os.environ.get('OPENSEARCH_USE_SSL', 1))) OPENSEARCH_TIMEOUT = 30 OPENSEARCH_MAX_RETRIES = 5 LOGGER_LEVEL = os.environ.get('LOGGER_LEVEL', INFO) LOGGER_DSN = os.environ.get('LOGGER_DSN') JSON_FILE_TEMPLATE = os.environ.get( 'JSON_FILE_TEMPLATE', 'yt_conflict_{status}.json') CSV_FILE_TEMPLATE = os.environ.get( 'CSV_FILE_TEMPLATE', 'yt_conflict_{status}.csv') S3_BUCKET = os.environ.get('S3_BUCKET', 'dev-orchdbucket') S3_FOLDER_TEMPLATE = os.environ.get( 'S3_FOLDER_TEMPLATE', 'yt_conflict_es_etl/{date}') ELASTICSEARCH_EXPORT_TEMP_PATH = 'tmp/elasticsearch_export/{date}' OS_CONFLICTS_INDEX_NAME = os.environ.get( 'OPENSEARCH_INDEX', 'conflict.v01') OS_CONFLICTS_ALIAS_NAME = os.environ.get( 'OPENSEARCH_ALIAS', 'conflict_write') CONFLICT_STATUSES = ('new',) CHUNK_SIZE = int(os.environ.get('CHUNK_SIZE', 500)) MAX_CHUNK_BYTES = int(os.environ.get('MAX_CHUNK_BYTES', 5242880)) THREAD_COUNT = int(os.environ.get('THREAD_COUNT', 4)) QUEUE_SIZE = int(os.environ.get('QUEUE_SIZE', 4)) SF_INSERT_BATCH_SIZE = int(os.environ.get('SF_INSERT_BATCH_SIZE', 15000)) SF_UPDATE_RESOLVED_BATCH_SIZE = int(os.environ.get( 'SF_UPDATE_RESOLVED_BATCH_SIZE', 10000))