【问题标题】:Prepare my bigdata with Spark via Python通过 Python 使用 Spark 准备我的大数据
【发布时间】:2017-01-17 00:30:27
【问题描述】:

我的 100m 大小的量化数据:

(1424411938', [3885, 7898])
(3333333333', [3885, 7898])

想要的结果:

(3885, [3333333333, 1424411938])
(7898, [3333333333, 1424411938])

所以我想要的是转换数据,以便我将 3885(例如)与所有拥有它的 data[0] 分组)。这是我在 中所做的:

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


【解决方案1】:

您可以使用一堆基本的 pyspark 转换来实现这一点。

>>> rdd = sc.parallelize([(1424411938, [3885, 7898]),(3333333333, [3885, 7898])])
>>> r = rdd.flatMap(lambda x: ((a,x[0]) for a in x[1]))

我们使用flatMapx[1] 中的每个项目设置一个键值对,并将数据行格式更改为(a, x[0]),这里的ax[1] 中的每个项目。要更好地理解flatMap,您可以查看文档。

>>> r2 = r.groupByKey().map(lambda x: (x[0],tuple(x[1])))

我们只是将所有键值对按键分组,并使用元组函数将可迭代对象转换为元组。

>>> r2.collect()
[(3885, (1424411938, 3333333333)), (7898, (1424411938, 3333333333))]

正如您所说,您可以使用 [:150] 来拥有前 150 个元素,我想这是正确的用法:

r2 = r.groupByKey().map(lambda x: (x[0],tuple(x[1])[:150]))

我尽量解释清楚。我希望这会有所帮助。

【讨论】:

  • 糟糕,我以为我评论过了。你能解释一下吗?例如,在我的实际应用程序中,我只想保留每个列表的 150 项,所以我想我可以对 tuple(x[1]) 进行切片。代码仍在运行,所以我不知道它是否可以解决问题...... :)
  • 谢谢,这也将帮助未来的用户,而不仅仅是我!我将等待代码执行然后采取行动。 :) 删除我的问题下的评论,不需要它(当你这样做时我也会编辑它):)
  • 我看到了与Active tasks is a negative number spark ui 相同的行为。我想知道问题出在代码上还是集群上发生了一些奇怪的事情......
  • 代码没问题,我的分区很多,我在那个链接中更新了答案,谢谢 paradosko。
猜你喜欢
  • 2016-02-06
  • 2017-11-08
  • 2017-11-24
  • 2017-03-02
  • 1970-01-01
  • 1970-01-01
  • 1970-01-01
  • 1970-01-01
  • 2021-02-22
相关资源
最近更新 更多