- c - 在位数组中找到第一个零
- linux - Unix 显示有关匹配两种模式之一的文件的信息
- 正则表达式替换多个文件
- linux - 隐藏来自 xtrace 的命令
我有一个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 /上下文:
session = SparkSession.builder.appName("myApp").getOrCreate()
session.stop()
根据 official documentation , getOrCreate
“获取一个现有的 SparkSession,或者,如果不存在,则根据此构建器中设置的选项创建一个新的。”但我不想要这种行为(这种行为是进程试图获取现有 session 的行为)。我找不到任何方法来禁用它,也不知道如何销毁 session ——我只知道如何停止其关联的 SparkContext。
如何在独立模块中构建新的 SparkSession,并在同一个 Python 进程中按顺序执行它们,而以前的 session 不会干扰新创建的 session ?
以下是项目结构的示例:
主.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
进程.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
最佳答案
长话短说,Spark(包括 PySpark)并非设计用于在单个应用程序中处理多个上下文。如果您对 JVM 方面的故事感兴趣,我建议您阅读 SPARK-2243 (解决为不会修复)。
PySpark 中做出了许多设计决策,包括但不限于 a singleton Py4J gateway .有效you cannot have multiple SparkContexts
in a single application . SparkSession
不仅绑定(bind)到 SparkContext
,还会引入其自身的问题,例如处理本地(独立)Hive metastore(如果使用的话)。此外,还有内部使用 SparkSession.builder.getOrCreate
的函数,取决于您现在看到的行为。一个值得注意的例子是 UDF 注册。如果存在多个 SQL 上下文,其他函数可能会表现出意外行为(例如 RDD.toDF
)。
多个上下文不仅不受支持,而且在我个人看来,还违反了单一职责原则。您的业务逻辑不应关注所有设置、清理和配置细节。
个人建议如下:
如果应用程序由多个连贯的模块组成,这些模块可以组合在一起并受益于具有缓存和通用元存储的单个执行环境,则在应用程序入口点初始化所有必需的上下文,并在必要时将这些上下文传递给各个管道:
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()
collect.py
/process.py
:
from pyspark.sql import SparkSession
def execute(spark: SparkSession, data=None):
...
否则(根据您的描述,这里似乎是这种情况)我会设计入口点来执行单个管道并使用外部工作流管理器(如 Apache Airflow 或 Toil )来处理执行。
它不仅更清洁,而且允许更灵活的故障恢复和调度。
同样的事情当然可以用构建器来完成,但就像 smart person曾经说过:显式优于隐式。
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()
collect.py
/process.py
与上一点相同。
无论哪种方式,我都会保留一个而且只有一个地方设置了上下文,并且只有一个地方被拆除。
关于python - 我怎样才能拆除一个 SparkSession 并在一个应用程序中创建一个新的?,我们在Stack Overflow上找到一个类似的问题: https://stackoverflow.com/questions/41491972/
前言: 有时候,一个数据库有多个帐号,包括数据库管理员,开发人员,运维支撑人员等,可能有很多帐号都有比较大的权限,例如DDL操作权限(创建,修改,删除存储过程,创建,修改,删除表等),账户多了,管理
所以我用 Create React App 创建并设置了一个大型 React 应用程序。最近我们开始使用 Storybook 来处理和创建组件。它很棒。但是,当我们尝试运行或构建应用程序时,我们不断遇
遵循我正在创建的控件的代码片段。这个控件用在不同的地方,变量也不同。 我正在尝试编写指令来清理代码,但在 {{}} 附近插入值时出现解析错误。 刚接触 Angular ,无法确定我错过了什么。请帮忙。
我正在尝试创建一个 image/jpeg jax-rs 提供程序类,它为我的基于 post rest 的 Web 服务创建一个图像。我无法制定请求来测试以下内容,最简单的测试方法是什么? @POST
我一直在 Windows 10 的模拟器中练习 c。后来我改用dev C++ IDE。当我在 C 中使用 FILE 时。创建的文件的名称为 test.txt ,而我给出了其他名称。请帮助解决它。 下面
当我们创建自定义 View 时,我们将 View 文件的所有者设置为自定义类,并使用 initWithFrame 或 initWithCode 对其进行实例化。 当我们创建 customUITable
我正在尝试为函数 * Producer 创建一个线程,但用于创建线程的行显示错误。我为这句话加了星标,但我无法弄清楚它出了什么问题...... #include #include #include
今天在做项目时,遇到了需要创建JavaScript对象的情况。所以Bing了一篇老外写的关于3种创建JavaScript对象的文章,看后跟着打了一遍代码。感觉方法挺好的,在这里与大家分享一下。 &
我正在阅读将查询字符串传递给 Amazon 的 S3 以进行身份验证的文档,但似乎无法理解 StringToSign 的创建和使用方式。我正在寻找一个具体示例来说明 (1) 如何构造 String
前言:我对 C# 中任务的底层实现不太了解,只了解它们的用法。为我在下面屠宰的任何东西道歉: 对于“我怎样才能开始一项任务但不等待它?”这个问题,我找不到一个好的答案。在 C# 中。更具体地说,即使任
我有一个由一些复杂的表达式生成的 ILookup。假设这是按姓氏查找人。 (在我们简单的世界模型中,姓氏在家庭中是唯一的) ILookup families; 现在我有两个对如何构建感兴趣的查询。 首
我试图创建一个 MSI,其中包含 和 exe。在 WIX 中使用了捆绑选项。这样做时出错。有人可以帮我解决这个问题。下面是代码: 错误 error LGH
在 Yii 中,Create 和 Update 通常使用相同的形式。因此,如果我在创建期间有电子邮件、密码、...other_fields...等字段,但我不想在更新期间专门显示电子邮件和密码字段,但
上周我一直在努力创建一个给定一行和一列的 QModelIndex。 或者,我会满足于在已经存在的 QModelIndex 中更改 row() 的值。 任何帮助,将不胜感激。 编辑: QModelInd
出于某种原因,这不起作用: const char * str_reset_command = "\r\nReset"; const char * str_config_command = "\r\nC
现在,我有以下由 original.df %.% group_by(Category) %.% tally() %.% arrange(desc(n)) 创建的 data.frame。 DF 5),
在今天之前,我使用/etc/vim/vimrc来配置我的vim设置。今天,我想到了创建.vimrc文件。所以,我用 touch .vimrc cat /etc/vim/vimrc > .vimrc 所
我可以创建一个 MKAnnotation,还是只读的?我有坐标,但我发现使用 setCooperative 手动创建 MKAnnotation 并不容易。 想法? 最佳答案 MKAnnotation
在以下代码中,第一个日志语句按预期显示小数,但第二个日志语句记录 NULL。我做错了什么? NSDictionary *entry = [[NSDictionary alloc] initWithOb
我正在使用与此类似的代码动态添加到数组; $arrayF[$f+1][$y][$x+1] = $value+1; 但是我在错误报告中收到了这个: undefined offset :1 问题:尝试创
我是一名优秀的程序员,十分优秀!