【问题标题】:how to specify the partition for mapPartition in spark如何在spark中为mapPartition指定分区
【发布时间】:2016-02-07 08:03:50
【问题描述】:

我想做的是分别计算每个列表,例如,如果我有 5 个列表 ([1,2,3,4,5,6],[2,3,4,5,6],[3,4,5,6],[4,5,6],[5,6]) 并且我想获得没有 6 个的 5 个列表,我会这样做:

data=[1,2,3,4,5,6]+[2,3,4,5,6,7]+[3,4,5,6,7,8]+[4,5,6,7,8,9]+[5,6,7,8,9,10]

def function_1(iter_listoflist):
    final_iterator=[]
    for sublist in iter_listoflist:
        final_iterator.append([x for x in sublist if x!=6])
    return iter(final_iterator)  

sc.parallelize(data,5).glom().mapPartitions(function_1).collect()

然后剪切列表,以便我再次获得第一个列表。 有没有办法简单地分离计算?我不希望这些列表混合在一起,它们的大小可能不同。

谢谢

菲利普

【问题讨论】:

  • 并不总是我只想并行计算列表的最后一个元素。并行化的条目是任何工作,这适用于相同大小的列表。如果有办法不使用并行化并直接给出分区,那就太好了。我只是希望它计算彼此分开的不同列表并给我不同的结果,这些结果也是列表

标签: python apache-spark pyspark partition


【解决方案1】:

据我了解您的意图,您只需在 parallelize 您的数据时将各个列表分开:

data = [[1,2,3,4,5,6], [2,3,4,5,6,7], [3,4,5,6,7,8],
    [4,5,6,7,8,9], [5,6,7,8,9,10]]

rdd = sc.parallelize(data)

rdd.take(1) # A single element of a RDD is a whole list
## [[1, 2, 3, 4, 5, 6]]

现在您可以使用您选择的函数简单地map

def drop_six(xs):
    return [x for x in xs if x != 6]

rdd.map(drop_six).take(3)
## [[1, 2, 3, 4, 5], [2, 3, 4, 5, 7], [3, 4, 5, 7, 8]]

【讨论】:

  • 我犯了一个错误并提供了错误的数据感谢您提供的一切帮助! @
  • 很高兴能帮上忙。
  • 实际上我只是尝试使用 mapPartitions 而不是 map,但它没有用!它给了我与输入相同的列表,我认为即使在stackoverflow.com/questions/21185092/… 和文档之后我也有一些不明白的地方
  • map 是否因为 parallelize 调用而并行计算?
  • 而且它不应该工作。此外,这里没有理由使用mapPartitionsmap over RDD 并行工作并简化了很多你可以认为这是parallelize的结果。
猜你喜欢
  • 2016-08-01
  • 2016-12-04
  • 2016-10-19
  • 2018-01-04
  • 1970-01-01
  • 2021-12-21
  • 2016-03-11
  • 1970-01-01
  • 1970-01-01
相关资源
最近更新 更多