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
相关产品推荐
相关产品推荐

