Snowflake性能测试:ETL管道改造后的无缓存测试与查询ID获取问题
1. 如何获取非生产环境下管道的查询ID?
在存储过程内记录查询ID
在存储过程执行核心Merge操作后,调用LAST_QUERY_ID()获取当前会话的最后执行查询ID,可将其写入自定义日志表或通过Airflow XCom返回:DECLARE v_query_id VARCHAR; BEGIN -- 你的Merge操作逻辑 MERGE INTO target_table t USING (SELECT * FROM your_view WHERE your_filter_conditions) s ON t.id = s.id WHEN MATCHED THEN UPDATE SET ... WHEN NOT MATCHED THEN INSERT (...); -- 获取并记录查询ID v_query_id := LAST_QUERY_ID(); INSERT INTO etl_query_logs (dag_id, query_id, exec_time) VALUES ('your_test_dag', v_query_id, CURRENT_TIMESTAMP()); END;之后可通过查询
etl_query_logs表获取目标ID。从Airflow任务日志提取
如果使用SnowflakeOperator调用存储过程,任务执行日志中通常会打印Snowflake返回的查询ID,直接进入Airflow UI对应任务实例的日志页面,搜索query_id或CALL等关键词即可找到。精准筛选Snowflake查询历史
通过INFORMATION_SCHEMA.QUERY_HISTORY视图,结合仓库名、执行时间、查询类型等条件缩小范围,注意嵌套查询需通过parent_query_id关联:SELECT query_id, query_text, start_time, parent_query_id FROM TABLE(INFORMATION_SCHEMA.QUERY_HISTORY( DATEADD('hour', -2, CURRENT_TIMESTAMP()), CURRENT_TIMESTAMP() )) WHERE warehouse_name = 'YOUR_TEST_WH' AND (query_text LIKE '%CALL your_stored_proc%' OR query_text LIKE '%MERGE INTO target_table%') ORDER BY start_time DESC;
2. 更便捷的性能测试方案
存储过程内预生成执行计划
无需实际执行Merge操作,直接在存储过程中用EXPLAIN生成执行计划并记录,避免缓存干扰:DECLARE v_plan_content VARCHAR; BEGIN -- 生成执行计划(仅解析不执行) SELECT LISTAGG(plan_line, '\n') INTO v_plan_content FROM TABLE(EXPLAIN( MERGE INTO target_table t USING (SELECT * FROM your_view WHERE your_filter_conditions) s ON t.id = s.id WHEN MATCHED THEN UPDATE SET ... WHEN NOT MATCHED THEN INSERT (...) )); -- 写入测试日志表 INSERT INTO etl_test_logs (test_scenario, log_content) VALUES ('merge_plan_with_new_filters', v_plan_content); END;模拟Airflow会话手动测试
在Snowflake客户端使用与Airflow相同的用户、仓库、会话参数(如ALTER SESSION SET USE_CACHED_RESULT=FALSE;),直接调用存储过程后立即执行:CALL your_stored_proc(); SELECT LAST_QUERY_ID(); SELECT * FROM TABLE(EXPLAIN_PLAN(LAST_QUERY_ID()));这种方式跳过Airflow调度,能快速验证执行计划和性能,适合迭代调试。
自动化性能对比脚本
用Snowflake Python连接器编写简单脚本,自动调用存储过程、捕获查询ID、拉取执行计划和耗时,方便多次测试对比:import snowflake.connector # 初始化连接 conn = snowflake.connector.connect( user="test_user", password="test_pass", account="your_account", warehouse="test_wh", database="test_db", schema="test_schema" ) cursor = conn.cursor() # 禁用缓存 cursor.execute("ALTER SESSION SET USE_CACHED_RESULT=FALSE;") # 调用存储过程 cursor.execute("CALL your_stored_proc();") # 获取查询ID并拉取详情 query_id = cursor.execute("SELECT LAST_QUERY_ID();").fetchone()[0] plan = cursor.execute(f"SELECT plan_line FROM TABLE(EXPLAIN_PLAN('{query_id}'));").fetchall() stats = cursor.execute(f"SELECT total_elapsed_time, bytes_scanned FROM TABLE(INFORMATION_SCHEMA.QUERY_HISTORY_BY_QUERY_ID('{query_id}'));").fetchone() # 输出结果 print(f"查询ID: {query_id}") print("执行计划:") for line in plan: print(line[0]) print(f"总耗时: {stats[0]}ms | 扫描字节数: {stats[1]}") cursor.close() conn.close()彻底清空缓存的测试环境
测试前先重启虚拟仓库清空数据缓存,同时设置会话级缓存参数:ALTER WAREHOUSE test_wh SUSPEND; ALTER WAREHOUSE test_wh RESUME; ALTER SESSION SET USE_CACHED_RESULT=FALSE; ALTER SESSION SET SESSION_CACHED_RESULT_TIMEOUT=0;确保每次测试都是冷启动状态,结果更准确。
内容的提问来源于stack exchange,提问作者Sam

