【发布时间】:2015-10-21 10:54:46
【问题描述】:
我有两个 RDD 说
rdd1 =
id | created | destroyed | price
1 | 1 | 2 | 10
2 | 1 | 5 | 11
3 | 2 | 3 | 11
4 | 3 | 4 | 12
5 | 3 | 5 | 11
rdd2 =
[1,2,3,4,5] # lets call these value as timestamps (ts)
rdd2 基本上是使用 range(initial_value, end_value, interval) 生成的。这里的参数可以变化。大小可以与 rdd1 相同或不同。这个想法是使用过滤条件根据 rdd2 的值将记录从 rdd1 提取到 rdd2(从 rdd1 的记录可以在提取时重复,如您在输出中看到的那样)
过滤条件 rdd1.created
预期输出:
ts | prices
1 | 10,11 # i.e. for ids 1,2 of rdd1
2 | 11,11 # ids 2,3
3 | 11,12,11 # ids 2,4,5
4 | 11,11 # ids 2,5
现在我想根据一些使用RDD2的键的条件过滤RDD1。 (如上所述)并返回连接RDD2的键和RDD1的过滤结果的结果
所以我这样做:
rdd2.map(lambda x : somefilterfunction(x, rdd1))
def somefilterfunction(x, rdd1):
filtered_rdd1 = rdd1.filter(rdd1[1] <= x).filter(rdd1[2] > x)
prices = filtered_rdd1.map(lambda x : x[3])
res = prices.collect()
return (x, list(res))
我得到:
例外:您似乎正在尝试广播 RDD 或 从动作或转换中引用 RDD。 RDD 转换 并且动作只能由驱动程序调用,不能在其他内部调用 转变;例如,rdd1.map(lambda x: rdd2.values.count() * x) 无效,因为值转换和计数操作 不能在 rdd1.map 转换内执行。更多 信息,请参阅 SPARK-5063。
我尝试使用 groupBy ,但由于这里 rdd1 的元素可以一次又一次地重复,而我知道分组只会将 rdd1 的每个元素组合在某个特定的插槽中一次。
现在唯一的方法是使用普通的 for 循环并进行过滤并最终加入所有内容。
有什么建议吗?
【问题讨论】:
标签: python pyspark apache-spark-sql rdd