Airflow技术问题:如何从BigQuery表获取数据并转为列表使用?
问题原因
你的代码核心问题在于Airflow的DAG解析时机和任务运行时机不匹配:
run_export这个Operator是在Airflow调度器加载DAG(解析阶段)时就初始化的,这时候process_data任务还没执行,bq_data还是初始的空列表,所以访问bq_data[0]会直接抛出"列表为空"的错误。process_data任务里的XCom数据只有在DAG实际运行时才会被拉取,这时候DAG已经完成解析,无法再修改run_export的参数。
解决方案
你可以通过两种方式实现动态生成CSV文件名的需求:
方案一:使用Jinja2模板直接引用XCom数据
Airflow支持用Jinja2模板在Operator参数中直接引用XCom数据,无需通过全局变量传递。注意任务ID必须在解析阶段确定,不能依赖运行时数据,所以任务ID要固定,只动态修改存储路径:
from airflow.providers.google.cloud.operators.bigquery import BigQueryGetDataOperator, BigQueryToCloudStorageOperator from airflow.operators.python import PythonOperator from airflow.models import DAG from datetime import datetime default_args = { 'start_date': datetime(2024, 1, 1) } with DAG('bq_export_dag', default_args=default_args, schedule_interval=None) as dag: get_data = BigQueryGetDataOperator( task_id='get_data_from_bq', dataset_id='SK22', table_id='current_times', max_results=100, # 此处传整数而非字符串 selected_fields='current_timestamps', ) def process_data_from_bq(**kwargs): ti = kwargs['ti'] bq_data = ti.xcom_pull(task_ids='get_data_from_bq') # 若需格式化时间戳等处理,可在此操作后推回XCom ti.xcom_push(key='formatted_timestamp', value=bq_data[0][0]) process_data = PythonOperator( task_id='process_data_from_bq', python_callable=process_data_from_bq, provide_context=True ) run_export = BigQueryToCloudStorageOperator( task_id="save_data_on_storage", # 任务ID固定,不可动态生成 source_project_dataset_table="a-data-set", # 用Jinja2模板引用处理后的XCom数据 destination_cloud_storage_uris=[ "gs://europe-west1-airflow-bucket/data/test{{ ti.xcom_pull(task_ids='process_data_from_bq', key='formatted_timestamp') }}.csv" ], export_format="CSV", field_delimiter=",", print_header=False, ) get_data >> process_data >> run_export
方案二:在PythonOperator中动态执行导出逻辑
如果需要更灵活的控制,可以直接在PythonOperator中使用BigQueryHook执行导出操作,在运行时拿到XCom数据后再构造导出路径:
from airflow.providers.google.cloud.hooks.bigquery import BigQueryHook from airflow.providers.google.cloud.operators.bigquery import BigQueryGetDataOperator from airflow.operators.python import PythonOperator from airflow.models import DAG from datetime import datetime default_args = { 'start_date': datetime(2024, 1, 1) } def export_to_gcs(**kwargs): ti = kwargs['ti'] # 拉取BigQuery返回的数据 bq_data = ti.xcom_pull(task_ids='get_data_from_bq') if not bq_data: raise ValueError("未从BigQuery获取到数据") timestamp = bq_data[0][0] # XCom返回格式为[(时间戳,), ...],需取嵌套元素 # 初始化Hook执行导出 hook = BigQueryHook(gcp_conn_id='google_cloud_default') hook.extract_table_to_storage( source_project_dataset_table="a-data-set", destination_cloud_storage_uris=[f"gs://europe-west1-airflow-bucket/data/test{timestamp}.csv"], export_format="CSV", field_delimiter=",", print_header=False ) with DAG('bq_export_dag', default_args=default_args, schedule_interval=None) as dag: get_data = BigQueryGetDataOperator( task_id='get_data_from_bq', dataset_id='SK22', table_id='current_times', max_results=100, selected_fields='current_timestamps', ) run_export = PythonOperator( task_id='save_data_on_storage', python_callable=export_to_gcs, provide_context=True ) get_data >> run_export
额外注意点
BigQueryGetDataOperator的max_results参数需传整数,而非字符串,否则可能触发异常。- XCom存储的数据格式为列表,每个元素是对应行的元组,因此访问首个值需要用
bq_data[0][0],而非bq_data[0]。
内容的提问来源于stack exchange,提问作者Empty Whale
相关产品推荐
相关产品推荐

