vae*_*r-k 12 python apache-spark pyspark
我有一个带有多个独立模块的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.
Run Code Online (Sandbox Code Playgroud)
这些模块中的每一个都在执行开始时构建SparkSession,并在其进程结束时停止sparkContext.我建立并停止会话/上下文,如下所示:
session = SparkSession.builder.appName("myApp").getOrCreate()
session.stop()
Run Code Online (Sandbox Code Playgroud)
根据官方文档,getOrCreate"获取现有的SparkSession,或者,如果没有现有的SparkSession,则根据此构建器中设置的选项创建一个新的SparkSession." 但我不希望这种行为(此过程尝试获取现有会话).我找不到任何方法来禁用它,我无法弄清楚如何破坏会话 - 我只知道如何停止其关联的SparkContext.
如何在独立模块中构建新的SparkSession,并在同一个Python进程中按顺序执行它们,而以前的会话不会干扰新创建的?
以下是项目结构的示例:
main.py
import collect
import process
if __name__ == '__main__':
data = collect.execute()
process.execute(data)
Run Code Online (Sandbox Code Playgroud)
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
Run Code Online (Sandbox Code Playgroud)
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
Run Code Online (Sandbox Code Playgroud)
简而言之,Spark(包括PySpark)并不是设计用于在单个应用程序中处理多个上下文.如果您对故事的JVM方面感兴趣,我建议您阅读SPARK-2243(已解决,因为无法修复).
在PySpark中做出了许多设计决策,这些决策反映了包括但不限于单独的Py4J网关.实际上,您不能SparkContexts在单个应用程序中拥有多个.SparkSession不仅受到约束,SparkContext而且还引入了自己的问题,如处理本地(独立)Hive Metastore(如果使用的话).此外,还有一些SparkSession.builder.getOrCreate内部使用的功能,取决于您现在看到的行为.一个值得注意的例子是UDF注册.如果存在多个SQL上下文,则其他函数可能会出现意外行为(例如RDD.toDF).
多个情境不仅不受支持,而且在我个人看来,也违反了单一责任原则.您的业务逻辑不应该关注所有设置,清理和配置细节.
我的个人建议如下:
如果应用程序由多个相关模块组成,这些模块可以组合并受益于具有缓存的单个执行环境,并且常见的Metastore会初始化应用程序入口点中的所有必需上下文,并在必要时将这些上下文传递给各个管道:
main.py:
from pyspark.sql import SparkSession
import collect
import process
if __name__ == "__main__":
spark: SparkSession = ...
# Pass data between modules
collected = collect.execute(spark)
processed = process.execute(spark, data=collected)
...
spark.stop()
Run Code Online (Sandbox Code Playgroud)collect.py/ process.py:
from pyspark.sql import SparkSession
def execute(spark: SparkSession, data=None):
...
Run Code Online (Sandbox Code Playgroud)否则(这似乎是基于你的描述的情况)我会设计入口点来执行单个管道并使用外部worfklow管理器(如Apache Airflow或Toil)来处理执行.
它不仅更清洁,而且还允许更灵活的故障恢复和调度.
同样的事情当然可以用建设者来完成,但就像一个聪明的人曾经说过的那样:明确比隐含更好.
main.py
import argparse
from pyspark.sql import SparkSession
import collect
import process
pipelines = {"collect": collect, "process": process}
if __name__ == "__main__":
parser = argparse.ArgumentParser()
parser.add_argument('--pipeline')
args = parser.parse_args()
spark: SparkSession = ...
# Execute a single pipeline only for side effects
pipelines[args.pipeline].execute(spark)
spark.stop()
Run Code Online (Sandbox Code Playgroud)collect.py/ process.py如上一点.
无论如何,我会保留一个且只有一个地方设置上下文,而且只有一个地方被拆除.
这是一种解决方法,但不是解决方案:
我发现源代码中的SparkSession类包含以下内容(我已从此处的显示中删除了不相关的代码行):__init__
_instantiatedContext = None
def __init__(self, sparkContext, jsparkSession=None):
self._sc = sparkContext
if SparkSession._instantiatedContext is None:
SparkSession._instantiatedContext = self
Run Code Online (Sandbox Code Playgroud)
因此,我可以通过设置解决方法我的问题_instantiatedContext在会议属性None调用后session.stop()。当下一个模块执行时,它调用getOrCreate()并没有找到前一个_instantiatedContext,因此它分配一个新的sparkContext.
这不是一个非常令人满意的解决方案,但它可以作为满足我当前需求的解决方法。我不确定启动独立会话的整个方法是反模式还是不寻常。
| 归档时间: |
|
| 查看次数: |
17097 次 |
| 最近记录: |