Airflow Composer执行BigQuery查询转DataFrame时conn_id未定义如何解决
问题根因及排查修复方向
核心错误是代码混淆了BigQueryHook的bigquery_conn_id参数与GCP项目ID的定义:bigquery_conn_id为Airflow后台预配置的连接资源名称,不是直接填写项目ID的参数,你当前代码在环境变量未配置GCP_PROJECT_ID时,会直接把默认值hard_coded_project_name当做连接ID查询,自然找不到对应连接。
具体排查步骤
- 修正BigQueryHook传参:Google Cloud Composer环境默认预置了名为
google_cloud_default的GCP连接,已经绑定了当前Composer环境的服务账号权限,直接把bigquery_conn_id参数设置为该值即可,不需要传入项目ID。 - 确认环境变量配置:如果需要自定义连接,先到Composer环境配置页检查是否已配置
GCP_PROJECT_ID环境变量,若未配置,代码会直接落到默认值hard_coded_project_name作为连接ID查找,导致报错。 - 权限校验:确认Composer环境的默认服务账号拥有BigQuery数据集
LP_RAW的查询权限,至少需要授予BigQuery Data Viewer、BigQuery Job User角色。
优化后代码示例
from airflow.models import DAG import os from airflow.operators.python_operator import PythonOperator import datetime from airflow.contrib.hooks.bigquery_hook import BigQueryHook default_args = { 'start_date': datetime.datetime(2020, 1, 1), } PROJECT_ID = os.environ.get("GCP_PROJECT_ID", "替换为你实际的GCP项目ID") def list_dates_in_df(): hook = BigQueryHook( bigquery_conn_id="google_cloud_default", use_legacy_sql=False, project_id=PROJECT_ID ) query = "select count(*) from LP_RAW.DIM_ACCOUNT;" # BigQueryHook自带转DataFrame方法,不需要手动初始化客户端 df = hook.get_pandas_df(query) # 后续可以加你自己的df处理逻辑 with DAG( 'df_test', schedule_interval=None, catchup = False, default_args=default_args ) as dag: list_dates = PythonOperator( task_id ='list_dates', python_callable = list_dates_in_df )
内容的提问来源于stack exchange,提问作者Lena Meer
相关产品推荐
相关产品推荐

