【问题标题】:How to utilize the same spark session in another module如何在另一个模块中使用相同的 spark 会话
【发布时间】: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


    【解决方案1】:

    这在两个不同的任务中是不可能的。每个 airlfow 任务都可能在不同的进程中运行,也可能在不同的机器上运行,因此当您执行任务时,您需要注意它们的运行彼此完全隔离,并且所有连接/会话等都由一个任务在内存中打开不会被带到第二个任务。

    任务之间的数据只能由 XComs 交换,这些数据要么存储在数据库中,要么(当您使用外部存储配置 Airflow 时)存储在外部存储(例如 S3/GCS)中。任务之间不共享内存中的数据。

    这是 Airflow 的非常基本的假设。

    如果您想在两个不同的步骤之间重用内存中的资源,您必须在 Airflow 中将它们设为当前的一个步骤。

    有计划在未来针对此类案件开放 - 但恕我直言,这还很遥远。

    【讨论】:

      猜你喜欢
      • 2012-12-18
      • 1970-01-01
      • 1970-01-01
      • 2019-09-09
      • 2019-02-19
      • 2012-10-13
      • 1970-01-01
      • 1970-01-01
      相关资源
      最近更新 更多