Python调用Ads Data Hub API批量查询及数据未写入BigQuery问题解决
Ads Data Hub API批量执行查询及数据未写入BigQuery问题解决
一、单个查询数据未写入BigQuery的排查与修复
关键错误点修复
name参数格式错误analysisQueries().start方法的name必须是查询的完整资源名称,格式为customers/{customer_id}/analysisQueries/{query_id},而非客户名称。你需要先通过list接口获取目标客户下的所有查询资源名,或者从ADH界面复制查询ID拼接。必填参数缺失/错误
adsDataCustomerId不能为空,填写对应的广告客户ID(纯数字,去掉横杠);destTable确保格式正确:推荐使用完整资源名projects/{your_project_id}/datasets/{dataset_id}/tables/{table_name},同时确认服务账号拥有该BigQuery数据集的Data Editor或Writer权限;startDate/endDate要和查询定义的参数匹配,如果查询本身定义了start_date和end_date参数,需要在spec.parameterValues里传递,而非直接写在spec根节点。
查询状态验证
执行start后返回的QueryMetadata中startTime为空,说明查询未成功启动。调用analysisQueries().get(name='查询资源名')获取实时状态,查看state字段(如RUNNING/SUCCEEDED/FAILED),如果是FAILED,查看error字段获取具体原因。
二、Python批量执行所有分析查询代码示例
from googleapiclient.discovery import build import time from google.oauth2.service_account import Credentials # 加载服务账号认证 credentials = Credentials.from_service_account_file('你的服务账号密钥文件路径.json') # 初始化ADH服务 service = build('adsdatahub', 'v1', credentials=credentials) # 配置核心参数 CUSTOMER_ID = '你的ADH客户ID(纯数字)' ADS_DATA_CUSTOMER_ID = '对应广告客户ID(纯数字)' PROJECT_ID = 'BigQuery项目ID' DATASET_ID = '目标数据集ID' TIME_ZONE = 'UTC' START_DATE = {'year': 2022, 'month': 12, 'day': 1} END_DATE = {'year': 2022, 'month': 12, 'day': 12} def list_all_analysis_queries(customer_id): """列出指定客户下的所有分析查询""" response = service.customers().analysisQueries().list( parent=f'customers/{customer_id}' ).execute() return response.get('analysisQueries', []) def start_adh_query(query_resource_name): """启动单个ADH查询,自动适配查询参数""" # 获取查询的参数定义 query_details = service.customers().analysisQueries().get( name=query_resource_name ).execute() # 构建参数值(根据查询实际参数调整) parameter_values = {} param_keys = [p['name'] for p in query_details.get('parameters', [])] if 'start_date' in param_keys: parameter_values['start_date'] = { 'value': f"{START_DATE['year']}-{START_DATE['month']:02d}-{START_DATE['day']:02d}" } if 'end_date' in param_keys: parameter_values['end_date'] = { 'value': f"{END_DATE['year']}-{END_DATE['month']:02d}-{END_DATE['day']:02d}" } if 'time_zone' in param_keys: parameter_values['time_zone'] = {'value': TIME_ZONE} # 生成唯一目标表名(避免重复) query_id = query_resource_name.split('/')[-1] dest_table = f'projects/{PROJECT_ID}/datasets/{DATASET_ID}/tables/adh_result_{query_id}' # 启动查询 return service.customers().analysisQueries().start( name=query_resource_name, body={ 'spec': { 'adsDataCustomerId': ADS_DATA_CUSTOMER_ID, 'parameterValues': parameter_values }, 'destTable': dest_table } ).execute() def wait_for_query_completion(query_run_name): """轮询等待查询执行完成""" while True: query_metadata = service.customers().analysisQueries().get( name=query_run_name ).execute() state = query_metadata.get('state') if state == 'SUCCEEDED': print(f"查询 {query_run_name} 执行成功,结果已写入: {query_metadata['destTable']}") return True elif state == 'FAILED': error_msg = query_metadata.get('error', {}).get('message', '未知错误') print(f"查询 {query_run_name} 执行失败: {error_msg}") return False else: print(f"查询 {query_run_name} 当前状态: {state},等待60秒后重试...") time.sleep(60) def main(): queries = list_all_analysis_queries(CUSTOMER_ID) if not queries: print("未找到任何可执行的分析查询") return print(f"共找到 {len(queries)} 个查询,开始批量执行...") for idx, query in enumerate(queries): query_name = query['name'] query_title = query['title'] print(f"\n=== 执行第 {idx+1}/{len(queries)} 个查询: {query_title} ===") try: start_response = start_adh_query(query_name) query_run_name = start_response['name'] wait_for_query_completion(query_run_name) except Exception as e: print(f"查询 {query_title} 执行出错: {str(e)}") continue if __name__ == '__main__': main()
核心注意事项
- 权限验证:服务账号需同时拥有ADH超级用户权限、BigQuery数据集写入权限(Data Editor角色)、项目BigQuery Job User权限;
- 参数适配:不同查询的自定义参数可能不同,需根据
list_all_analysis_queries返回的parameters字段动态调整; - 并发控制:ADH有查询并发限制,若查询数量较多,可添加并发控制逻辑(如每次仅运行2-3个查询);
- 错误重试:可对失败查询添加重试机制,避免因临时网络问题导致任务终止。
内容的提问来源于stack exchange,提问作者SKP
相关产品推荐
相关产品推荐

