Airflow中S3ListOperator如何返回数据?仅靠XCom或PythonOperator吗?
问题解答
1. S3ListOperator 获取文件列表并推送XCom的方法
你当前的代码里,S3ListOperator默认不会把文件列表推送到XCom,所以xcom_pull返回None。只需添加do_xcom_push=True参数(Airflow 2.x版本),就能让它把扫描到的文件列表推送到XCom:
response_table_parquet_list = S3ListOperator( task_id='get_bucket_object_list_operator', aws_conn_id='cloudFlare_conn', bucket='paprika-table-storage', prefix='response_table/', do_xcom_push=True, # 开启XCom推送,将文件列表传入XCom dag=dag, )
如果是Airflow 1.x版本,参数名是xcom_push,替换成xcom_push=True即可。
之后在后续任务中,就能通过ti.xcom_pull(task_ids='get_bucket_object_list_operator')拿到完整的文件列表了。
2. 并非只有PythonOperator能返回数据/使用XCom
很多Airflow内置Operator都支持XCom,只要Operator内部实现了XCom推送逻辑,或者提供了开关参数(比如上面的do_xcom_push):
- BashOperator:可以通过
xcom_push=True捕获命令输出推送到XCom; - BigQueryOperator:可以设置
do_xcom_push=True将查询结果推送到XCom; - S3GetObjectOperator:支持将获取的文件内容推送到XCom;
这些Operator不需要依赖PythonOperator就能完成数据传递,核心看Operator本身是否支持XCom推送能力。
内容的提问来源于stack exchange,提问作者lima
相关产品推荐
相关产品推荐

