【发布时间】:2015-08-21 11:50:16
【问题描述】:
我有一个用 scala 编写的 Akka 系统,它需要调用一些 Python 代码,依赖于 Pandas 和 Numpy,所以我不能只使用 Jython。我注意到 Spark 在其工作节点上使用 CPython,所以我很好奇它如何执行 Python 代码以及该代码是否以某种可重用的形式存在。
【问题讨论】:
标签: scala pandas apache-spark interop pyspark
我有一个用 scala 编写的 Akka 系统,它需要调用一些 Python 代码,依赖于 Pandas 和 Numpy,所以我不能只使用 Jython。我注意到 Spark 在其工作节点上使用 CPython,所以我很好奇它如何执行 Python 代码以及该代码是否以某种可重用的形式存在。
【问题讨论】:
标签: scala pandas apache-spark interop pyspark
此处描述了 PySpark 架构 https://cwiki.apache.org/confluence/display/SPARK/PySpark+Internals。
正如@Holden 所说,Spark 使用 py4j 从 python 访问 JVM 中的 Java 对象。但这只是一种情况——当驱动程序是用 python 编写时(图的左侧)
另一种情况(图右侧)——Spark Worker 启动 Python 进程,将序列化的 Java 对象发送给 Python 程序进行处理,并接收输出。 Java 对象被序列化为 pickle 格式 - 因此 python 可以读取它们。
看起来您正在寻找的是后一种情况。这里有一些指向 Spark 的 scala 核心的链接,可能对您入门很有用:
Pyrolite 库,为 Python 的 pickle 协议提供 Java 接口 - Spark 使用该库将 Java 对象序列化为 pickle 格式。例如,在访问 PairRDD 的 Key、Value 对的 Key 部分时需要进行这种转换。
启动 python 进程并对其进行迭代的 Scala 代码:api/python/PythonRDD.scala
选择代码的 SerDeser 实用程序:api/python/SerDeUtil.scala
Python 端:python/pyspark/worker.py
【讨论】:
所以 Spark 使用 py4j 在 JVM 和 Python 之间进行通信。这允许 Spark 使用不同版本的 Python,但需要序列化来自 JVM 的数据,反之亦然以进行通信。 http://py4j.sourceforge.net/ 有更多关于 py4j 的信息,希望对您有所帮助:)
【讨论】: