import os import random from timeit import default_timer as timer import numpy as np from snowflake import connector import json from contextlib import contextmanager sf_config = { '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'), 'password': os.environ.get('SNOWFLAKE_PASSWORD'), 'account': os.environ.get('SNOWFLAKE_ACCOUNT') } def get_connection(): return connector.connect( user=sf_config['user'], password=sf_config['password'], account=sf_config['account'], warehouse=sf_config['warehouse'], database=sf_config['db'], schema=sf_config['schema'], role=sf_config['role'], autocommit=True) @contextmanager def get_cursor(cn): cursor = cn.cursor() try: yield cursor finally: cursor.close() def query(conn=None): def _query(c): with get_cursor(c) as cursor: session_id = c.session_id count = 0 sql = """ SELECT fa.releaseId AS release_id, cat.catalogName AS catalog_name, da.artistName AS artist_name, fa.catalogId AS catalog_id, fa.artistId AS artist_id, fa.activities FROM (SELECT MIN(artistId) AS artistId, MIN(releaseId) as releaseId, catalogId, SUM(units) AS activities FROM facts.prod.fact_analytics WHERE labelId = 21786 AND dayId BETWEEN 7117 AND 7147 AND transactionTypeId IN(12, 23) AND storeId NOT IN(3, 334, 399, 446, 312, 36) GROUP BY catalogId) AS fa INNER JOIN facts.prod.dim_artist da ON fa.artistId = da.artistId INNER JOIN facts.prod.dim_catalog cat ON cat.catalogId = fa.catalogId ORDER BY activities DESC LIMIT 10 """.format() # params = {'labelid': random.choice(label_ids), # 'isrc': random.choice(isrcs), 'limit': limit, } cursor.execute(sql) for _ in cursor: count += 1 return session_id, cursor.sfqid, count if not conn: conn = get_connection() result = _query(conn) conn.close() return result else: return _query(conn) def get_sys_time(conn, session_id, query_id): with get_cursor(conn) as cursor: cursor.execute(""" SELECT total_elapsed_time FROM TABLE(information_schema.query_history_by_session( session_id=>%(session_id)s) ) WHERE query_id = %(query_id)s AND execution_status = 'SUCCESS'; """, {'session_id': session_id, 'query_id': query_id}) sys_time, *_ = cursor.fetchone() return float(sys_time) / 1000.0 def calc_metrics(l): return { 'avg': np.average(l), 'median': np.median(l), '95p': np.percentile(l, 95), 'max': max(l) } def analyze(n=13): print('WH: {}'.format(sf_config['warehouse'])) usr_times = [] sys_times = [] dif_times = [] conn = get_connection() for i in range(n): start = timer() session_id, query_id, row_cnt = query(conn) usr_time = timer() - start sys_time = get_sys_time(conn, session_id, query_id) print('{}: (usr: {:.3f}; sys: {:.3f}) [Rows: {}]'.format( session_id, usr_time, sys_time, row_cnt)) usr_times.append(usr_time) sys_times.append(sys_time) dif_times.append(usr_time - sys_time) conn.close() print('Usr time: {}'.format(json.dumps(calc_metrics(usr_times), indent=4))) print('Sys time: {}'.format(json.dumps(calc_metrics(sys_times), indent=4))) print('Dif time: {}'.format(json.dumps(calc_metrics(dif_times), indent=4))) def main(): analyze() if __name__ == '__main__': main()