Python如何循环遍历多个BQ表批量完成数据打分处理
问题根因
Jupyter环境的
%%bigquery是单元格级魔法命令,仅在单元格顶层生效,无法在for循环的缩进块内运行,也不能自动解析循环中的动态表名变量,这是魔法命令本身的执行机制限制,和循环逻辑、缩进语法无关。itertools是迭代器工具,不解决魔法命令的执行域问题,手动逐表运行效率极低。
解决方案
直接使用BigQuery官方Python客户端替代单元格魔法命令,支持在循环内动态传入表名查询,后续已验证正常的单表处理、预测逻辑不需要做任何改动,即可实现批量自动跑批。
前置依赖
如果环境未安装BQ客户端,先执行命令安装:pip install google-cloud-bigquery[pandas]
安装后和原有%%bigquery使用的权限、配置完全通用,不需要额外做认证调整。
修正后可直接运行的代码
import pandas as pd import numpy as np from google.cloud import bigquery # 初始化BigQuery客户端 client = bigquery.Client() # BQ分表列表 scoring_tables = [ "`Customer_Analytics.HS_AF_PROPERTY_ALL_DATA_SCORE_PV_CLEANED_01`", "`Customer_Analytics.HS_AF_PROPERTY_ALL_DATA_SCORE_PV_CLEANED_02`", "`Customer_Analytics.HS_AF_PROPERTY_ALL_DATA_SCORE_PV_CLEANED_03`", "`Customer_Analytics.HS_AF_PROPERTY_ALL_DATA_SCORE_PV_CLEANED_04`", "`Customer_Analytics.HS_AF_PROPERTY_ALL_DATA_SCORE_PV_CLEANED_05`" ] # 存储所有预测结果 all_scores = [] for t in scoring_tables: # 动态拼接SQL查询分表,结果直接转pandas DataFrame,替代原%%bigquery魔法逻辑 query_sql = f"SELECT * FROM {t}" score_data = client.query(query_sql).to_dataframe() # 以下为已验证可正常运行的单表处理逻辑,无需修改 score_data.set_index('HH_ID', inplace=True) # HOUSE_INCOME、AGE字段零值填充 score_data['HOUSE_INCOME'] = np.where(score_data['HOUSE_INCOME'] == 0, 107, score_data['HOUSE_INCOME']) score_data['AGE'] = np.where(score_data['AGE'] == 0, 54, score_data['AGE']) # 房屋属性字段重分类 # 外墙类型重分类 wall_cond = [ score_data['PROP_EXTR_WALL_TYPE'].str.contains("BRICK"), score_data['PROP_EXTR_WALL_TYPE'].str.contains("WOOD"), score_data['PROP_EXTR_WALL_TYPE'].str.contains("CONCRETE"), score_data['PROP_EXTR_WALL_TYPE'].str.contains("METAL"), score_data['PROP_EXTR_WALL_TYPE'].str.contains("STEEL") ] wall_choice = ["BRICK", "WOOD", "CONCRETE", "METAL", "METAL"] score_data['PROP_EXTR_WALL_TYPE_MOD'] = np.select(wall_cond, wall_choice, default="OTHERS") # 车库类型重分类 grg_cond = [ score_data['PROP_GRG_TYPE'].str.contains("ATTACHED"), score_data['PROP_GRG_TYPE'].str.contains("DETACHED"), score_data['PROP_GRG_TYPE'].str.contains("CARPORT"), score_data['PROP_GRG_TYPE'].str.contains("BASEMENT") ] grg_choice = ["ATTACHED", "DETACHED", "CARPORT", "BASEMENT"] score_data['PROP_GRG_TYPE_MOD'] = np.select(grg_cond, grg_choice, default="OTHERS") # 屋顶类型重分类 roof_cond = [ score_data['PROP_ROOF_TYPE'].str.contains("GABLE"), score_data['PROP_ROOF_TYPE'].str.contains("HIP"), score_data['PROP_ROOF_TYPE'].str.contains("GAMBREL") ] roof_choice = ["GABLE", "HIP", "GAMBREL"] score_data['PROP_ROOF_TYPE_MOD'] = np.select(roof_cond, roof_choice, default="OTHERS") # 指定字段独热编码 to_encode = ["IND_ETHNICITY","IND_GENDER","IND_MOVERS_FLAG","IND_OCCUPATION","IND_REGION","PROP_EXTR_WALL_TYPE","PROP_GRG_TYPE","PROP_ROOF_TYPE","TSI"] score_data_dm = pd.get_dummies(data=score_data, columns=to_encode, drop_first=False) # 对齐训练集字段结构 missing_cols = [c for c in train_X.columns if c not in score_data_dm.columns] score_data_dm[missing_cols] = 0 score_data_dm_filt = score_data_dm[select_columns] # 模型预测,结果追加到总列表 y_pred = xgb_prop_PV.predict_proba(score_data_dm_filt)[:, 1] all_scores.extend(y_pred)
优化点说明
- 原代码中
all_scores = all_scores + y_pred改为all_scores.extend(y_pred),避免每次循环新建列表对象,跑批效率更高 - 所有数据处理、预测逻辑和手动单表运行的逻辑完全一致,不会出现结果偏差
- 后续如果新增分表,只需要把表名加到
scoring_tables列表即可,不需要修改其他代码
内容的提问来源于stack exchange,提问作者Piyush
相关产品推荐
相关产品推荐

