Airflow 2.5.3下SparkSubmitOperator从Spark应用推送XCom的可行性及方案咨询
SparkSubmitOperator 推送数据到XCom的问题解答
关键信息
- Airflow 2.5.3 版本的
SparkSubmitOperator确实没有内置在Spark应用内部直接推送数据到XCom的功能 - 截至当前的最新官方版本,该算子仍未添加此特性
- 实现需求主要有两种途径:自定义算子扩展,或采用间接传递方案
具体实现方式
自定义算子扩展
基于原SparkSubmitOperator进行改造即可,核心思路:
- 让Spark应用将需要传递的数据写入中间存储(比如本地临时文件、Redis、数据库等)
- 自定义算子继承
SparkSubmitOperator,重写execute方法,在Spark任务执行完成后读取中间存储中的数据 - 调用
task_instance.xcom_push()方法将数据推送至XCom
间接传递(无需自定义算子)
如果不想修改算子,可以拆分任务流程:
- 第一个任务使用原
SparkSubmitOperator执行Spark应用,将需要传递的数据写入外部存储(如HDFS、关系型数据库) - 第二个任务通过PythonOperator或其他合适的算子,读取外部存储的数据后推送至XCom,供后续任务使用
内容的提问来源于stack exchange,提问作者sou_joshi
相关产品推荐
相关产品推荐

