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

如何在Airflow中从Big Query获取SQL查询结果并转为pandas DataFrame?

在Airflow中从BigQuery获取结果并转为Pandas DataFrame的完整实现

你的代码框架已经很靠谱了!我来帮你补全所有细节,让这个DAG能顺利跑起来:

1. 先确保依赖安装到位

首先要确认你的Airflow环境已经安装了Google云提供商的插件包,这是使用BigQueryHook的前提:

pip install apache-airflow-providers-google

2. 补全Python函数的完整代码

你需要导入必要的模块,还可以按需添加XCom传递(如果要把DataFrame传给后续任务的话),完整函数如下:

import logging
from airflow.providers.google.cloud.hooks.bigquery import BigQueryHook

def get_my_query(*args, **kwargs):
    # 初始化BigQuery Hook,使用你配置好的连接ID
    bq_hook = BigQueryHook(bigquery_conn_id='my_connection_id', delegate_to=None, use_legacy_sql=False)
    
    # 定义你的BigQuery查询
    my_query = """
        SELECT col1, col2 
        FROM `my_bq_project.my_bq_dataset.my_table`
    """
    
    # 执行查询并转为Pandas DataFrame
    df = bq_hook.get_pandas_df(sql=my_query)  # 注意:新版BigQueryHook参数名已从bql改为sql,旧版本可保留bql
    logging.info(f"查询结果前5行:\n{df.head()}")
    
    # 如果需要把DataFrame传给后续任务,用XCom推送(大数据量不建议,优先存到存储服务)
    kwargs['ti'].xcom_push(key='bq_query_result', value=df.to_json(orient='split'))
    
    return df

3. 补全DAG的完整定义

把你的DAG片段补全,添加PythonOperator来执行上面的函数:

from airflow import DAG
from airflow.operators.python import PythonOperator
from datetime import datetime, timedelta

# DAG默认参数配置
default_args = {
    'owner': 'airflow',
    'depends_on_past': False,
    'start_date': datetime(2024, 1, 1),
    'email_on_failure': False,
    'email_on_retry': False,
    'retries': 1,
    'retry_delay': timedelta(minutes=5),
}

# 定义完整DAG
with DAG(
    'my_bq_to_df_dag',
    default_args=default_args,
    description='从BigQuery查询数据并转为Pandas DataFrame',
    schedule_interval='@daily',  # 按需调整调度频率,比如@hourly或者具体 cron 表达式
    catchup=False,
    tags=['bigquery', 'pandas'],
) as my_dag:

    # 定义Python任务
    fetch_bq_data = PythonOperator(
        task_id='fetch_bq_data',
        python_callable=get_my_query,
        provide_context=True,  # 必须设为True才能让函数获取kwargs中的任务实例(ti)
        dag=my_dag,
    )

    # 如果有后续任务,可在这里添加依赖关系
    # fetch_bq_data >> next_task

4. 关键注意事项

  • Airflow连接配置:你需要在Airflow UI的「Admin > Connections」里创建ID为my_connection_id的BigQuery连接,可以选择用服务账号密钥配置,或者在GCP部署的Airflow直接使用工作负载身份。
  • 权限设置:连接BigQuery的身份(服务账号或默认身份)必须拥有目标表的roles/bigquery.dataViewer权限,以及项目的roles/bigquery.jobUser权限,否则会出现权限报错。
  • 大数据量处理:如果你的查询结果数据量很大,不要用XCom传递(有大小限制),建议直接把DataFrame写入GCS、本地存储或者其他数据库。

5. 测试方法

你可以在Airflow UI里手动触发这个DAG,然后查看任务日志,就能看到打印出来的DataFrame前5行。如果用了XCom,也可以在任务实例的「XCom」标签页里查看推送的JSON格式数据。

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

相关产品推荐
方舟 Agent Plan

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

最近更新时间:2026.05.21 08:44:01