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

Python调用BigQuery执行大查询遇内存错误的优化问询

解决BigQuery查询内存溢出问题的优化方案

首先,咱们核心的问题就是把大数据量的查询结果加载到本地DataFrame——这完全没必要!BigQuery本身就支持直接把查询结果写入目标表,跳过本地内存这一步,这是解决内存错误的关键。下面分步骤给你优化:

一、最关键的优化:跳过DataFrame,直接让BigQuery写入目标表

原来的流程是 BigQuery查询 → 本地DataFrame → 写入BigQuery,现在改成 BigQuery查询 → 直接写入目标表,彻底绕开本地内存的限制。具体做法是在QueryJobConfig里直接指定目标表,不需要再调用to_dataframe()。

二、优化BigQuery查询语句,减少计算和扫描量

你的原始查询有几个可以精简的点,能大幅降低数据处理量:

  • 合并重复CTE:Q3和Q4都是查询有交易用户的首次/末次订单日期,合并成一个CTE能减少一次全表扫描
  • 移除冗余UNNEST:Q1里的CROSS JOIN UNNEST(hits.product)根本没用到对应字段,会生成大量冗余数据,直接删掉
  • 避免重复计算:把多次用到的session_id计算提前放到基础CTE里,减少重复运算
  • 统一过滤逻辑:把重复的用户维度过滤条件合并,避免冗余判断

优化后的查询语句如下:

DECLARE target_date DATE;
SET target_date = DATE_SUB(CURRENT_DATE(), INTERVAL @days_ago day);

WITH user_base AS (
    SELECT 
        customDimension.value AS UserID,
        CONCAT(CAST(fullVisitorId AS STRING),CAST(visitId AS STRING)) AS session_id,
        hits.transaction.transactionId AS transaction_id,
        Date AS visit_date,
        totals.bounces AS bounces,
        totals.transactionRevenue AS raw_revenue
    FROM `project.dataset.ga_sessions_20*` AS t
    CROSS JOIN UNNEST (hits) AS hits
    CROSS JOIN UNNEST(t.customdimensions) AS customDimension
    WHERE 
        PARSE_DATE('%y%m%d', _table_suffix) BETWEEN DATE_SUB(target_date, INTERVAL 999 day) AND target_date
        AND customDimension.index = 2
        AND customDimension.value NOT IN ("true", "false", "undefined")
        AND customDimension.value IS NOT NULL
),
user_session_metrics AS (
    SELECT 
        UserID,
        COUNT(DISTINCT session_id) AS visits,
        COUNT(DISTINCT transaction_id) AS orders,
        SAFE_DIVIDE(COUNT(DISTINCT transaction_id), COUNT(DISTINCT session_id)) AS conversion_rate,
        MIN(visit_date) AS first_visit_date,
        MAX(visit_date) AS last_visit_date,
        IFNULL(SUM(bounces), 0) AS bounces,
        IFNULL(SUM(raw_revenue)/1000000, 0) AS revenue
    FROM user_base
    GROUP BY UserID
),
user_order_dates AS (
    SELECT 
        customDimension.value AS UserID,
        MIN(Date) AS first_order,
        MAX(Date) AS last_order
    FROM `project.dataset.ga_sessions_*` AS t
    CROSS JOIN UNNEST(t.customdimensions) AS customDimension
    WHERE 
        totals.transactions > 0
        AND customDimension.index = 2
        AND customDimension.value NOT IN ("true", "false", "undefined")
        AND customDimension.value IS NOT NULL
    GROUP BY UserID
)
SELECT 
    a.UserID,
    a.visits,
    IFNULL(a.orders, 0) AS orders,
    IFNULL(a.conversion_rate, 0) AS conversion_rate,
    IFNULL(a.bounces, 0) AS bounces,
    IFNULL(a.revenue, 0) AS revenue,
    IFNULL(SAFE_DIVIDE(a.revenue, a.orders), 0) AS AOV,
    IFNULL(SAFE_DIVIDE(a.revenue, a.visits), 0) AS rev_per_visit,
    IFNULL(SAFE_DIVIDE(a.bounces, a.visits), 0) AS bounce_rate
FROM user_session_metrics AS a
LEFT JOIN user_order_dates AS b USING (UserID)
GROUP BY a.UserID, a.visits, a.orders, a.conversion_rate, a.bounces, a.revenue

三、修改Python函数,适配新的流程

现在不需要再生成DataFrame,直接让BigQuery把查询结果写入目标表,同时用参数计算目标日期,不用从DataFrame里提取:

def generate_user_pit(days_ago):
    ### 生成指定天数前的用户事实表并保存至BigQuery(优化版)
    client = bigquery.Client()  # 若全局已初始化client,可省略此行
    
    # 计算目标日期,用于生成表名
    target_date_query = client.query(f"SELECT DATE_SUB(CURRENT_DATE(), INTERVAL {days_ago} day)")
    target_date = target_date_query.result().to_dataframe().iloc[0,0]
    table_name = f"data_{target_date}"
    dataset_ref = client.dataset('test')
    table_ref = dataset_ref.table(table_name)
    
    # 加载优化后的查询语句
    query = """
        DECLARE target_date DATE;
        SET target_date = DATE_SUB(CURRENT_DATE(), INTERVAL @days_ago day);

        WITH user_base AS (
            SELECT 
                customDimension.value AS UserID,
                CONCAT(CAST(fullVisitorId AS STRING),CAST(visitId AS STRING)) AS session_id,
                hits.transaction.transactionId AS transaction_id,
                Date AS visit_date,
                totals.bounces AS bounces,
                totals.transactionRevenue AS raw_revenue
            FROM `project.dataset.ga_sessions_20*` AS t
            CROSS JOIN UNNEST (hits) AS hits
            CROSS JOIN UNNEST(t.customdimensions) AS customDimension
            WHERE 
                PARSE_DATE('%y%m%d', _table_suffix) BETWEEN DATE_SUB(target_date, INTERVAL 999 day) AND target_date
                AND customDimension.index = 2
                AND customDimension.value NOT IN ("true", "false", "undefined")
                AND customDimension.value IS NOT NULL
        ),
        user_session_metrics AS (
            SELECT 
                UserID,
                COUNT(DISTINCT session_id) AS visits,
                COUNT(DISTINCT transaction_id) AS orders,
                SAFE_DIVIDE(COUNT(DISTINCT transaction_id), COUNT(DISTINCT session_id)) AS conversion_rate,
                MIN(visit_date) AS first_visit_date,
                MAX(visit_date) AS last_visit_date,
                IFNULL(SUM(bounces), 0) AS bounces,
                IFNULL(SUM(raw_revenue)/1000000, 0) AS revenue
            FROM user_base
            GROUP BY UserID
        ),
        user_order_dates AS (
            SELECT 
                customDimension.value AS UserID,
                MIN(Date) AS first_order,
                MAX(Date) AS last_order
            FROM `project.dataset.ga_sessions_*` AS t
            CROSS JOIN UNNEST(t.customdimensions) AS customDimension
            WHERE 
                totals.transactions > 0
                AND customDimension.index = 2
                AND customDimension.value NOT IN ("true", "false", "undefined")
                AND customDimension.value IS NOT NULL
            GROUP BY UserID
        )
        SELECT 
            a.UserID,
            a.visits,
            IFNULL(a.orders, 0) AS orders,
            IFNULL(a.conversion_rate, 0) AS conversion_rate,
            IFNULL(a.bounces, 0) AS bounces,
            IFNULL(a.revenue, 0) AS revenue,
            IFNULL(SAFE_DIVIDE(a.revenue, a.orders), 0) AS AOV,
            IFNULL(SAFE_DIVIDE(a.revenue, a.visits), 0) AS rev_per_visit,
            IFNULL(SAFE_DIVIDE(a.bounces, a.visits), 0) AS bounce_rate
        FROM user_session_metrics AS a
        LEFT JOIN user_order_dates AS b USING (UserID)
        GROUP BY a.UserID, a.visits, a.orders, a.conversion_rate, a.bounces, a.revenue
    """
    
    query_params = [
        bigquery.ScalarQueryParameter('days_ago', 'INT64', days_ago),
    ]
    
    job_config = bigquery.QueryJobConfig()
    job_config.query_parameters = query_params
    job_config.write_disposition = 'WRITE_TRUNCATE'
    # 指定目标表,直接将查询结果写入
    job_config.destination = table_ref
    
    print(f"Starting query to populate {table_name}...")
    query_job = client.query(query, job_config=job_config)
    # 等待任务执行完成
    query_job.result()
    print(f"Successfully created {table_name} in BigQuery!")

额外的优化小贴士

  • 开启查询缓存:如果查询逻辑不变、数据源未更新,BigQuery会直接返回缓存结果,节省时间和成本
  • 设置批量优先级:若任务非紧急,可添加job_config.priority = bigquery.QueryPriority.BATCH,降低查询成本
  • 监控任务状态:通过query_job.state和query_job.total_bytes_processed可查看任务进度和扫描数据量,方便排查问题

内容的提问来源于stack exchange,提问作者Ben P

相关产品推荐
方舟 Agent Plan

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

最近更新时间:2026.05.12 03:59:57