我正在尝试以下代码,它为RDD中的每一行添加一个数字,并使用PySpark返回一个RDD列表.
from pyspark.context import SparkContext
file = "file:///home/sree/code/scrap/sample.txt"
sc = SparkContext('local', 'TestApp')
data = sc.textFile(file)
splits = [data.map(lambda p : int(p) + i) for i in range(4)]
print splits[0].collect()
print splits[1].collect()
print splits[2].collect()
Run Code Online (Sandbox Code Playgroud)
输入文件(sample.txt)中的内容是:
1
2
3
Run Code Online (Sandbox Code Playgroud)
我期待这样的输出(在rdd中分别添加0,1,2的数字):
[1,2,3]
[2,3,4]
[3,4,5]
Run Code Online (Sandbox Code Playgroud)
而实际产出是:
[4, 5, 6]
[4, 5, 6]
[4, 5, 6]
Run Code Online (Sandbox Code Playgroud)
这意味着理解仅使用变量i的值3,而不考虑范围(4).
为什么会出现这种情况?
我想多进程一个列表 RDDS的如下
from pyspark.context import SparkContext
from multiprocessing import Pool
def square(rdd_list):
def _square(i):
return i*i
return rdd_list.map(_square)
sc = SparkContext('local', 'Data_Split')
data = sc.parallelize([1,2,3,4,5,6])
dataCollection = [data, data, data]
p = Pool(processes=2)
result = p.map(square, dataCollection)
print result[0].collect()
Run Code Online (Sandbox Code Playgroud)
我期待输出中的RDD列表,每个元素包含来自数据的平方元素.
但是运行代码会导致以下错误:
例外:您似乎正在尝试广播RDD或从动作或转换中引用RDD.RDD转换和操作只能由驱动程序调用,而不能在其他转换内部调用; 例如,rdd1.map(lambda x:rdd2.values.coun\t()*x)无效,因为无法在rdd1.map转换中执行值转换和计数操作.有关更多信息,请参阅SPARK-5063.
我的问题是: -
1)为什么代码不能按预期工作?我怎样才能解决这个问题 ?
2)如果我在我的RDD列表中使用p.map(池)而不是简单的映射,我是否会获得我的程序的任何性能增强(在减少运行时方面).