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)
产量
你好约翰
你好
你好乔治
你好保罗
您可以通过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)