【发布时间】:2021-11-23 08:18:00
【问题描述】:
我必须在 Airflow 中运行两个模块并执行两个任务。每个任务都有一个 PySpark 模块,用于执行一些 spark 操作。第二个模块使用在前一个会话中创建的数据框并继续其操作。
我们怎样才能用同样的SparkSession 初始化来达到同样的效果?我曾尝试使用getActiveSession(),但由于任务 1 作业已完成,它不起作用,因此当任务 2 运行时,会创建一个新的 spark 会话。
- [root@ ..dags]# cat tmp_spark_1.py
from pyspark.sql import SparkSession
spark = SparkSession.builder.appName("PRJT").enableHiveSupport().getOrCreate()
df = spark.createDataFrame([(1, 2), (3, 4)], ['a', 'b'])
df.show()
- [root@ ..dags]# cat tmp_spark_2.py
from pyspark.sql import SparkSession
#spark = SparkSession.builder.appName("PRJT").enableHiveSupport().getActiveSession().getOrCreate()
spark = SparkSession.getActiveSession()
df1 = df.select(df['a'])
df1.show()
【问题讨论】:
标签: python apache-spark pyspark apache-spark-sql airflow