import logging import os import datetime import snowflake.connector from snowflake.connector import DictCursor from datadog import initialize, api import boto3 from botocore.exceptions import ClientError import json logger = logging.getLogger() logger.setLevel(logging.INFO) def get_secret(): secret = "" secret_name = "dev/snowflake-slow-queries-monitor/db-creds" endpoint_url = "https://secretsmanager.us-east-1.amazonaws.com" region_name = "us-east-1" session = boto3.session.Session() client = session.client( service_name='secretsmanager', region_name=region_name, endpoint_url=endpoint_url ) try: get_secret_value_response = client.get_secret_value( SecretId=secret_name ) logger.info("Retrieved secret from secrets manager") except ClientError as e: if e.response['Error']['Code'] == 'ResourceNotFoundException': print("The requested secret " + secret_name + " was not found") elif e.response['Error']['Code'] == 'InvalidRequestException': print("The request was invalid due to:", e) elif e.response['Error']['Code'] == 'InvalidParameterException': print("The request had invalid params:", e) elif 'Error' in e.response: print("Error:", e) else: # Decrypted secret using the associated KMS CMK # Depending on whether the secret was a string or binary, one of these fields will be populated if 'SecretString' in get_secret_value_response: secret = get_secret_value_response['SecretString'] else: secret = get_secret_value_response['SecretBinary'] return secret def handler(event, context): creds = json.loads(get_secret()) logger.info("Get creds from secret manager for user: %s" % creds['SNOWFLAKE_USER']) user = creds['SNOWFLAKE_USER'] password = creds['SNOWFLAKE_PASSWORD'] account = os.environ.get('SNOWFLAKE_ACCOUNT') database = os.environ.get('SNOWFLAKE_DATABASE') schema = os.environ.get('SNOWFLAKE_SCHEMA') warehouse = os.environ.get('SNOWFLAKE_WAREHOUSE') timezone = os.environ.get('SNOWFLAKE_TIMEZONE') datadog_options = { 'api_key': creds['DD_API_KEY'], 'app_key': creds['DD_APP_KEY'] } initialize(**datadog_options) logger.info("Connecting to DB: %s" % database) try: snowflake_context = snowflake.connector.connect( account=account, user=user, password=password, database=database, schema=schema, warehouse=warehouse, timezone=timezone ) logger.info("Successfully connected to Snowflake") # Querying data by DictCursor cur = snowflake_context.cursor(DictCursor) try: cur.execute("select * \ from table(orchard_app_reporting.information_schema.query_history()) \ where \ EXECUTION_TIME > 3600000 \ AND EXECUTION_STATUS = 'RUNNING' \ order by start_time;") for rec in cur: print('{0}, {1}'.format(rec['WAREHOUSE_NAME'], rec['EXECUTION_TIME'])) ev_user_name = rec['USER_NAME'] logger.info("USER_NAME: %s" % ev_user_name) ev_warehouse_name = rec['WAREHOUSE_NAME'] logger.info("WAREHOUSE_NAME: %s" % ev_warehouse_name) ev_execution_time = datetime.datetime.fromtimestamp(rec['EXECUTION_TIME']).strftime('%H:%M:%S') logger.info("EXECUTION_TIME: %s" % ev_execution_time) ev_execution_minutes = (rec['EXECUTION_TIME']/1000/60) logger.info("EXECUTION_TIME_IN_MINS: %s" % ev_execution_minutes) title = "Query found running for over an hour on Snowflake" text = "User %s has been running a query on %s for %s minutes" % (ev_user_name, ev_warehouse_name, ev_execution_minutes) logger.info("Event: %s" % text) tags = ['warehouse:%s' % ev_warehouse_name, 'application:datawarehouse', 'username:%s' % ev_user_name] api.Event.create(title=title, text=text, tags=tags) logger.info("Event sent to Datadog") finally: cur.close() return {"message": 1} except Exception as ex: logger.error(ex) return {"message": 0}