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

Python调用Ads Data Hub API批量查询及数据未写入BigQuery问题解决

Ads Data Hub API批量执行查询及数据未写入BigQuery问题解决

一、单个查询数据未写入BigQuery的排查与修复

关键错误点修复

  1. name参数格式错误
    analysisQueries().start方法的name必须是查询的完整资源名称,格式为customers/{customer_id}/analysisQueries/{query_id},而非客户名称。你需要先通过list接口获取目标客户下的所有查询资源名,或者从ADH界面复制查询ID拼接。

  2. 必填参数缺失/错误

    • 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根节点。
  3. 查询状态验证
    执行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

相关产品推荐
方舟 Agent Plan

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

最近更新时间:2026.08.09 01:25:18