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

Airflow实操:如何将REST API数据加载至BigQuery?

Airflow从REST API导数据到BigQuery的修正方案

原代码存在的问题

  • PythonOperator指定的python_callable名称错误(写了download_sales但实际函数是insert_from_api_to_bq)
  • 仅实现了API数据打印,缺少核心的BigQuery插入逻辑
  • 代码未闭合(PythonOperator定义缺少结尾括号)
  • 未定义headers和default_args等必要变量
  • 使用urlopen不够健壮,建议用requests库简化API请求

修正后的完整代码

首先确保安装必要依赖:

pip install requests google-cloud-bigquery apache-airflow
from airflow import DAG
from airflow.operators.python import PythonOperator
from datetime import datetime
import requests
from google.cloud import bigquery
import json

# 定义DAG默认参数
default_args = {
    'owner': 'airflow',
    'retries': 1,
}

def insert_from_api_to_bq():
    # 根据API要求配置请求头,比如添加认证信息
    headers = {
        'Content-Type': 'application/json',
        # 示例认证头:'Authorization': 'Bearer YOUR_API_TOKEN'
    }

    # 发送API请求并处理响应
    response = requests.get('[URL GOES HERE]', headers=headers)
    response.raise_for_status()  # 自动抛出HTTP请求错误
    sale_data = response.json()

    # 初始化BigQuery客户端
    bq_client = bigquery.Client()
    # 替换为你的BigQuery表ID:项目名.数据集名.表名
    target_table = 'your-gcp-project.your-dataset.sales_table'

    # 整理数据格式,匹配BigQuery表结构
    rows_to_insert = []
    for sale in sale_data['SaleList']:
        print(f"sale ID: {sale['SaleID']} Customer:{sale['Customer']} Order Date: {sale['OrderDate']}")
        # 构造符合表结构的行数据
        row = {
            'SaleID': sale['SaleID'],
            'Customer': sale['Customer'],
            'OrderDate': sale['OrderDate']
            # 如有其他字段,按需添加
        }
        rows_to_insert.append(row)

    # 批量插入数据到BigQuery
    if rows_to_insert:
        insert_errors = bq_client.insert_rows_json(target_table, rows_to_insert)
        if not insert_errors:
            print("数据成功写入BigQuery")
        else:
            print(f"写入失败,错误信息:{insert_errors}")

with DAG(
    "sales_data_pipeline",
    start_date=datetime(2021, 1, 1),
    schedule_interval="@daily",
    default_args=default_args,
    catchup=False
) as dag:

    extract_load_sales = PythonOperator(
        task_id="extract_load_sales",
        python_callable=insert_from_api_to_bq
    )

关键注意事项

  • 确保Airflow运行账号拥有BigQuery表的写入权限
  • 替换[URL GOES HERE]为实际的REST API地址
  • 替换target_table为你的BigQuery项目、数据集和目标表名
  • 调整headers中的认证信息以匹配API要求
  • 确保BigQuery表的字段与代码中row字典的键完全匹配(包括字段名大小写)

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

相关产品推荐
方舟 Agent Plan

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

最近更新时间:2026.08.09 21:10:25