【发布时间】:2017-01-17 00:30:27
【问题描述】:
我的 100m 大小的量化数据:
(1424411938', [3885, 7898])
(3333333333', [3885, 7898])
想要的结果:
(3885, [3333333333, 1424411938])
(7898, [3333333333, 1424411938])
所以我想要的是转换数据,以便我将 3885(例如)与所有拥有它的 data[0] 分组)。这是我在python 中所做的:
def prepare(data):
result = []
for point_id, cluster in data:
for index, c in enumerate(cluster):
found = 0
for res in result:
if c == res[0]:
found = 1
if(found == 0):
result.append((c, []))
for res in result:
if c == res[0]:
res[1].append(point_id)
return result
但是当我 mapPartitions()'ed data RDD 和 prepare() 时,它似乎只在当前分区中执行我想要的操作,因此返回的结果比预期的要大。
例如,如果开始的第一个记录在第一个分区中,第二个在第二个分区,那么我会得到这样的结果:
(3885, [3333333333])
(7898, [3333333333])
(3885, [1424411938])
(7898, [1424411938])
如何修改我的prepare() 以获得想要的效果?或者,如何处理prepare()产生的结果,才能得到想要的结果?
正如您可能已经从代码中注意到的那样,我根本不在乎速度。
这是一种创建数据的方法:
data = []
from random import randint
for i in xrange(0, 10):
data.append((randint(0, 100000000), (randint(0, 16000), randint(0, 16000))))
data = sc.parallelize(data)
【问题讨论】:
-
是
DataFrame吗? -
没有@AlbertoBonsanto。
-
如果是 DataFrame 会更容易
-
我对 DF @AlbertoBonsanto 不是很熟悉,也许我应该得到... ;)
标签: python python algorithm apache-spark distributed-computing bigdata