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

使用Python云函数将HTTP API数据导入BigQuery遇阻求助

问题排查与解决方案

核心问题解析

你的代码在数据转DataFrame阶段出现致命错误,导致后续无法正常写入BigQuery:

  • 你将JSON数组反复序列化/反序列化为NDJSON字符串后,用pd.DataFrame.from_records(ndjson)创建DataFrame——这个方法接收的是字典/列表对象,传入字符串会被拆解为单个字符的行,完全不符合预期结构,直接引发写入BigQuery时的类型不匹配错误。
  • 冗余的JSON操作:new_json_response2已经是API返回的字典列表,无需多次json.dumps和json.loads。
  • 重复导入模块:函数内部重复导入requests、json等,属于无效代码,且影响可读性。

修正后的代码

import pandas as pd
import requests
from requests.structures import CaseInsensitiveDict
import pandas_gbq


def validate_http(request):
    # 无论请求类型,直接调用数据拉取函数(可根据需求调整逻辑)
    get_api_data()
    return f'Data pull complete'


def get_api_data():
    # 获取AccessToken
    headers = {'Content-Type': 'application/x-www-form-urlencoded'}
    data = f'client_id={my_client_id}&client_secret={my_client_secret}&grant_type=client_credentials&scope={my_scope}'
    
    token_response = requests.post(
        'https://login.microsoftonline.com/4fa9c138-d3e7-4bc3-8bab-a74bde6b7584/oauth2/v2.0/token',
        headers=headers,
        data=data
    )
    token_response.raise_for_status()  # 捕获请求失败的情况
    access_token = token_response.json()["access_token"]

    # 拉取业务数据
    url = "my_api_url"
    headers = CaseInsensitiveDict()
    headers["Accept"] = "application/json"
    headers["Authorization"] = f"Bearer {access_token}"
    
    api_resp = requests.get(url, headers=headers)
    api_resp.raise_for_status()  # 捕获API请求失败
    api_data = api_resp.json()["value"]

    # 直接用字典列表创建DataFrame,无需多余JSON操作
    df = pd.DataFrame.from_records(api_data)

    # 写入BigQuery
    bq_load('TABLE_NAME', df)


def bq_load(table_name, df):
    project_name = 'PROJECT_ID'
    dataset_name = 'DATASET_NAME'
    
    df.to_gbq(
        destination_table=f'{dataset_name}.{table_name}',
        project_id=project_name,
        if_exists='replace',
        progress_bar=False  # Cloud Function环境建议关闭进度条
    )

额外注意事项

  1. 权限配置:确保Cloud Function使用的服务账号拥有以下BigQuery权限:
    • bigquery.dataEditor:用于读写数据表
    • bigquery.jobUser:用于提交加载任务
  2. 数据类型匹配:BigQuery表的字段类型需与DataFrame的列类型严格对应(比如字符串对应STRING,数字对应INT64/FLOAT64等)
  3. 错误捕获:添加raise_for_status()可以及时发现API请求失败的情况,便于排查问题

替代方案:使用BigQuery客户端库直接导入(无需Pandas)

如果数据量较大,直接用BigQuery原生客户端库更高效,避免Pandas的内存开销:

from google.cloud import bigquery
import requests
from requests.structures import CaseInsensitiveDict


def validate_http(request):
    get_api_data()
    return f'Data pull complete'


def get_api_data():
    # 获取AccessToken(同前,省略重复代码)
    # ...
    api_resp = requests.get(url, headers=headers)
    api_resp.raise_for_status()
    api_data = api_resp.json()["value"]

    # 初始化BigQuery客户端
    client = bigquery.Client(project='PROJECT_ID')
    table_ref = client.dataset('DATASET_NAME').table('TABLE_NAME')

    # 直接写入数据
    errors = client.insert_rows_json(table_ref, api_data)
    if errors:
        raise RuntimeError(f'BigQuery写入失败: {errors}')

内容的提问来源于stack exchange,提问作者Selvin Kumar

相关产品推荐
方舟 Agent Plan

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

最近更新时间:2026.08.16 05:55:47