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

基于Spark与Dataproc优化BigQuery查询响应时间的技术咨询

Spark读取BigQuery性能优化方案

问题描述

使用Spark读取BigQuery中3400万行、30列(10个字符串、5个整数、其余为浮点数)、大小6GB的表,集群为20台N1-4核/16GB虚拟机的Dataproc集群。存在以下性能问题:

  • 查询有匹配结果时响应时间<1秒,但无匹配结果时耗时长达30-60秒,增加集群节点无明显提升,无法满足快速fallback需求。
  • 加载DataFrame并打印速度较快,但计算均值步骤耗时极长。

相关代码

from pyspark.sql import SparkSession
from pyspark.sql.functions import mean
spark = SparkSession.builder \
  .appName('1.2. BigQuery Storage & Spark SQL - Python')\
  .config('spark.jars', 'gs://spark-lib/bigquery/spark-bigquery-latest_2.12.jar') \
  .getOrCreate()


spark.conf.set("spark.sql.repl.eagerEval.enabled",True)

table = "pokerservers.parquetdb.DB19022023_flop"
df_query = spark.read \
  .format("bigquery") \
  .option("table", table) \
  .load()

df_query.createOrReplaceTempView("df_query_view")

df_query_view = spark.sql("""
FROM df_query_view
WHERE best_hand_flop = 'ONE_PAIR'
AND prob >= 0.9
AND prob <= 1.1
AND proba_win_flop >= 0.87
AND proba_win_flop >= 0.91
AND proba_draw_flop >= -0.01
AND proba_draw_flop <= 0.01
AND actions_preflop = 'rre'
AND actions_flop = 'e' 
AND actions_turn = '-'
AND actions_river = '-' 
AND checker_flop >= -0.01
AND checker_flop <= 0.01
AND Pair = 0 
AND Three = 0
AND Four = 0
AND Straight = 1
AND Colors = 1
AND Highcard = 13
AND Low = 1
AND Bottom = 1 
AND Middle = 0
AND Top = 1
AND StraightDirect = 1 
""").cache()

df_query_view.cache()
avg_call = df_query_view.select(mean('Call')).collect()[0][0]
print(avg_call)

优化建议

  • 下推过滤逻辑到BigQuery端:当前代码是全表加载到Spark后再过滤,无匹配时需扫描完整6GB数据。改为将WHERE条件通过filter参数传递给BigQuery数据源,让BigQuery先完成过滤,只返回符合条件的行(无匹配时直接返回空结果),避免全表扫描。示例:
    df_query = spark.read \
      .format("bigquery") \
      .option("table", table) \
      .option("filter", """
        best_hand_flop = 'ONE_PAIR'
        AND prob >= 0.9 AND prob <= 1.1
        AND proba_win_flop >= 0.91
        AND proba_draw_flop >= -0.01 AND proba_draw_flop <= 0.01
        AND actions_preflop = 'rre' AND actions_flop = 'e' 
        AND actions_turn = '-' AND actions_river = '-' 
        AND checker_flop >= -0.01 AND checker_flop <= 0.01
        AND Pair = 0 AND Three = 0 AND Four = 0
        AND Straight = 1 AND Colors = 1 AND Highcard = 13
        AND Low = 1 AND Bottom = 1 AND Middle = 0
        AND Top = 1 AND StraightDirect = 1
      """) \
      .load()
    
  • 移除冗余缓存和优化缓存时机:代码中连续两次调用df_query_view.cache()属于重复操作,删除其中一次。同时确保缓存操作在过滤完成后执行,避免缓存全表数据。
  • 清理冗余查询条件:WHERE子句中proba_win_flop >= 0.87被proba_win_flop >= 0.91覆盖,直接删除前者,减少不必要的条件判断。
  • 直接在BigQuery端计算均值:如果仅需要符合条件的Call列均值,无需将数据拉取到Spark计算。通过Spark读取BigQuery的查询结果,让BigQuery完成过滤和聚合:
    query_result = spark.read \
      .format("bigquery") \
      .option("query", """
        SELECT mean(Call) as avg_call
        FROM pokerservers.parquetdb.DB19022023_flop
        WHERE best_hand_flop = 'ONE_PAIR'
        AND prob >= 0.9 AND prob <= 1.1
        AND proba_win_flop >= 0.91
        AND proba_draw_flop >= -0.01 AND proba_draw_flop <= 0.01
        AND actions_preflop = 'rre' AND actions_flop = 'e' 
        AND actions_turn = '-' AND actions_river = '-' 
        AND checker_flop >= -0.01 AND checker_flop <= 0.01
        AND Pair = 0 AND Three = 0 AND Four = 0
        AND Straight = 1 AND Colors = 1 AND Highcard = 13
        AND Low = 1 AND Bottom = 1 AND Middle = 0
        AND Top = 1 AND StraightDirect = 1
      """) \
      .load()
    avg_call = query_result.collect()[0][0]
    
  • 调整Spark读取配置:启用BigQuery Storage API的ARROW格式传输,提升数据读取效率,添加配置:
    .option("useStorageApi", "true") \
    .option("readDataFormat", "ARROW")
    
  • 关闭eagerEval:spark.sql.repl.eagerEval.enabled开启后会自动触发DataFrame的action(如显示前几行),额外消耗资源,建议在生产代码中关闭该配置。
  • 排查集群资源瓶颈:通过Spark UI查看Stage执行详情,确认无匹配时的耗时是否集中在数据读取阶段,检查节点的IO、CPU利用率,排查是否存在数据倾斜或Shuffle阻塞问题。

内容的提问来源于stack exchange,提问作者Raphaël Ambit

相关产品推荐
方舟 Agent Plan

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

最近更新时间:2026.07.28 15:55:08