You need to enable JavaScript to run this app.
优惠活动
大模型
产品
解决方案
定价
更多

Snowflake性能测试:ETL管道改造后的无缓存测试与查询ID获取问题

针对Airflow调度的Snowflake ETL管道性能测试问题的解决方案

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

相关产品推荐
方舟 Agent Plan

超全模态模型 × Harness 升级,最新支持 Deepseek-V4.1-Flash、GLM-5.3 系列、Doubao-Seedream-5.0-pro、Kimi-K3 (部分), 限时 9.9 元起

最近更新时间:2026.08.05 21:50:33