Airflow中动态指定GCS源文件至BigQuery导入任务的实现方案
问题
我需要通过PythonOperator将格式为account_{year}_{month}.csv的文件(例如本月是account_2023_06.csv,下月是account_2023_07.csv)上传至Google Cloud Storage(GCS),后续再用GoogleCloudStorageToBigQueryOperator将该文件导入BigQuery表。但由于文件名每月动态变化,无法固定source_objects参数的取值。尝试过移除文件名中的年月标识(不符合需求),或者直接取存储桶中最新修改的文件(逻辑过于简单),请问该如何解决?
解决方案:通过Airflow XCom传递动态文件名
利用Airflow的XCom组件在任务间传递动态生成的文件路径,完美适配这种文件名随时间变化的场景,具体实现步骤如下:
1. 在PythonOperator的调用函数中推送文件路径到XCom
在生成并上传文件的Python函数里,通过任务实例(ti)的xcom_push方法,将生成的GCS文件路径存入XCom:
def accounts(**kwargs): auth = (kwargs['username'], kwargs['password']) ti = kwargs['ti'] ... # 初始化GCS客户端 client = storage.Client() bucket = client.get_bucket(BUCKET_accounts_historical_sl) # 生成文件路径并上传至GCS path = "/accounts/account_" + str(yr) + "_" + str(mnth) + ".csv" ti.xcom_push(key='path', value=path)
2. 为PythonOperator开启上下文传递
创建PythonOperator任务时,添加provide_context=True参数,确保函数能获取到Airflow的任务实例等上下文信息:
download_dataset_task = PythonOperator( task_id="download_dataset_task", python_callable=accounts, provide_context=True, )
3. 在GoogleCloudStorageToBigQueryOperator中拉取XCom的值
使用Airflow的模板语法,通过ti.xcom_pull从XCom中读取之前存储的文件路径,作为source_objects的参数值:
dataset_CS_BQ = GoogleCloudStorageToBigQueryOperator( task_id="dataset_CS_BQ", bucket=BUCKET_accounts_historical, source_objects=["{{ ti.xcom_pull(key='path') }}"], # 补充其他必要参数,比如目标表信息、文件格式等 )
内容的提问来源于stack exchange,提问作者Ekaterina Ponkratova
相关产品推荐
相关产品推荐

