Airflow多任务场景下如何跨PySpark模块复用同一个SparkSession
Airflow的每个任务默认是独立调度的进程,甚至可能运行在不同的worker节点上,第一个任务执行结束后进程直接销毁,其内存中维护的SparkSession、DataFrame变量也会一并释放,第二个任务启动是全新的进程空间,自然无法通过getActiveSession()获取到前序的会话。除此之外,第一个任务内存中定义的df变量也无法被第二个独立进程直接访问,哪怕拿到同一个SparkSession也不能直接调用df.select。
方案1:合并为单个Airflow任务(成本最低)
如果两个PySpark模块没有必须拆分调度的强需求,最简单的方式是把两个模块的逻辑合并到同一个Airflow任务的执行逻辑中,同一进程内天然共享同一个SparkSession,也可以直接传递DataFrame变量。
示例代码:
# 合并后的任务执行逻辑 from pyspark.sql import SparkSession # 导入两个模块封装好的业务函数 from tmp_spark_1 import run_module1 from tmp_spark_2 import run_module2 if __name__ == "__main__": spark = SparkSession.builder.appName("PRJT").enableHiveSupport().getOrCreate() # 运行第一个模块返回计算好的df df = run_module1(spark) # 直接把df传给第二个模块使用 run_module2(spark, df) spark.stop()
你需要把原来两个脚本里的执行逻辑封装成可调用的函数,不要把业务逻辑直接写在脚本顶层。
方案2:通过共享存储传递DataFrame(生产环境最常用)
如果必须拆分为两个独立的Airflow任务,采用「中间持久化」的方式传递状态是工业界通用方案:
- 第一个任务执行完成后,将生成的DataFrame写入共享存储(HDFS、对象存储、或者Hive临时表)
- 第二个任务启动后从共享存储读取数据生成新的DataFrame继续处理
示例代码:
# tmp_spark_1.py 末尾添加持久化逻辑 df.write.mode("overwrite").saveAsTable("tmp.tmp_prjt_calculated_df") # 也可以存为parquet文件:df.write.parquet("hdfs://path/to/tmp_df.parquet", mode="overwrite")
# tmp_spark_2.py 调整逻辑 from pyspark.sql import SparkSession spark = SparkSession.builder.appName("PRJT").enableHiveSupport().getOrCreate() # 读取临时表生成df df = spark.table("tmp.tmp_prjt_calculated_df") # 读parquet的方式:df = spark.read.parquet("hdfs://path/to/tmp_df.parquet") df1 = df.select(df['a']) df1.show()
如果不需要长期留存中间数据,第二个任务执行完成后可以删除对应的临时表/临时文件。
方案3:通过独立Spark服务复用会话(复杂度较高)
如果确实需要跨任务复用同一个SparkSession,可以先启动一个独立的长生命周期Spark服务(比如Spark Thrift Server、Livy服务),两个Airflow任务都作为客户端连接该服务提交计算逻辑,SparkSession由服务端统一维护,就可以跨请求共享会话和临时表数据。该方案需要额外维护Spark服务的可用性、配置会话超时策略,适合有大量跨任务Spark计算需求的场景。
内容的提问来源于stack exchange,提问作者Rocky1989

