如何在同一个Spark项目中同时使用Scala和Python?

Wil*_*iao 15 python scala apache-spark spark-streaming pyspark

是否可以将Spark RDD传递给Python?

因为我需要一个python库来对我的数据进行一些计算,但我的主要Spark项目是基于Scala的.有没有办法混合它们或让python访问相同的火花上下文?

小智 19

您确实可以使用Scala和Spark以及常规Python脚本来管理python脚本.

test.py

#!/usr/bin/python

import sys

for line in sys.stdin:
  print "hello " + line
Run Code Online (Sandbox Code Playgroud)

火花壳(斯卡拉)

val data = List("john","paul","george","ringo")

val dataRDD = sc.makeRDD(data)

val scriptPath = "./test.py"

val pipeRDD = dataRDD.pipe(scriptPath)

pipeRDD.foreach(println)
Run Code Online (Sandbox Code Playgroud)

产量

你好约翰

你好

你好乔治

你好保罗

  • 是的,我知道这个方法,但是python脚本在执行器上运行,所以我有一个问题,如果我将过多的数据输送到外部脚本,工作者会崩溃吗?我的意思是,外部Python脚本不是并行计算. (3认同)

Aja*_*pta 7

您可以通过Spark中的Pipe运行Python代码.

使用pipe(),您可以编写RDD的转换,将标准输入中的每个RDD元素作为String读取,根据脚本指令操作该String,然后将结果作为String写入标准输出.

SparkContext.addFile(路径),我们可以添加的文件列表中的每个工作节点的当星火工作starts.All工作节点都会有自己的脚本,因此我们将通过管道越来越并行操作的复制下载.我们需要在所有worker和executor节点上安装所有库和依赖项.

示例:

Python文件:将输入数据设置为大写的代码

#!/usr/bin/python
import sys
for line in sys.stdin:
    print line.upper()
Run Code Online (Sandbox Code Playgroud)

Spark代码:用于管道数据

val conf = new SparkConf().setAppName("Pipe")
val sc = new SparkContext(conf)
val distScript = "/path/on/driver/PipeScript.py"
val distScriptName = "PipeScript.py"
sc.addFile(distScript)
val ipData = sc.parallelize(List("asd","xyz","zxcz","sdfsfd","Ssdfd","Sdfsf"))
val opData = ipData.pipe(SparkFiles.get(distScriptName))
opData.foreach(println)
Run Code Online (Sandbox Code Playgroud)