【发布时间】:2017-01-05 18:18:17
【问题描述】:
我有一个带有多个独立模块的 pyspark 程序,每个模块可以独立处理数据以满足我的各种需求。但它们也可以链接在一起以处理管道中的数据。这些模块中的每一个都构建了一个 SparkSession 并自行完美执行。
但是,当我尝试在同一个 python 进程中连续运行它们时,我遇到了问题。在管道中的第二个模块执行的那一刻,spark 抱怨我尝试使用的 SparkContext 已停止:
py4j.protocol.Py4JJavaError: An error occurred while calling o149.parquet.
: java.lang.IllegalStateException: Cannot call methods on a stopped SparkContext.
这些模块中的每一个都在执行开始时构建一个 SparkSession,并在其进程结束时停止 sparkContext。我像这样构建和停止会话/上下文:
session = SparkSession.builder.appName("myApp").getOrCreate()
session.stop()
根据official documentation,getOrCreate“获取现有的 SparkSession,或者,如果没有现有的,则根据此构建器中设置的选项创建一个新的。”但我不想要这种行为(进程试图获取现有会话的这种行为)。我找不到任何方法来禁用它,也无法弄清楚如何销毁会话——我只知道如何停止其关联的 SparkContext。
如何在独立的模块中构建新的 SparkSession,并在同一个 Python 进程中按顺序执行它们,而以前的会话不会干扰新创建的会话?
以下是项目结构示例:
main.py
import collect
import process
if __name__ == '__main__':
data = collect.execute()
process.execute(data)
collect.py
import datagetter
def execute(data=None):
session = SparkSession.builder.appName("myApp").getOrCreate()
data = data if data else datagetter.get()
rdd = session.sparkContext.parallelize(data)
[... do some work here ...]
result = rdd.collect()
session.stop()
return result
process.py
import datagetter
def execute(data=None):
session = SparkSession.builder.appName("myApp").getOrCreate()
data = data if data else datagetter.get()
rdd = session.sparkContext.parallelize(data)
[... do some work here ...]
result = rdd.collect()
session.stop()
return result
【问题讨论】:
标签: python apache-spark pyspark