【问题标题】:Spark with own map and reduce functions pythonSpark具有自己的map和reduce函数python
【发布时间】:2015-08-05 00:43:52
【问题描述】:

我正在尝试使用 python spark 进行类似 mapreduce 的操作。这是我所拥有的和我的问题。

object_list = list(objects) #this is precomputed earlier in my script
def my_map(obj):
    return [f(obj)]
def my_reduce(obj_list1, obj_list2):
    return obj_list1 + obj_list2

我想要做的事情如下:

myrdd = rdd(object_list) #objects are now spread out
myrdd.map(my_map)
myrdd.reduce(my_reduce)
my_result = myrdd.result()

my_result 现在应该是 = [f(obj1), f(obj2), ..., f(objn)]。我想纯粹为了速度而使用 spark,我的脚本在 forloop 中执行此操作时花费了很长时间。有谁知道如何在 spark 中执行上述操作?

【问题讨论】:

    标签: python mapreduce apache-spark


    【解决方案1】:

    它通常看起来像这样:

    myrdd = sc.parallelize(object_list)
    my_result = myrdd.map(f).reduce(lambda a,b:a+b)
    

    有一个用于 RDD 的 sum 函数,所以这也可能是:

    myrdd = sc.parallelize(object_list)
    my_result = myrdd.map(f).sum()
    

    但是,这将为您提供一个数字。 f(obj1)+f(obj2)+...

    如果您想要一个包含所有响应 [f(obj1),f(obj2), ...] 的数组,则不要使用 .reduce().sum(),而是使用 .collect()

    myrdd = sc.parallelize(object_list)
    my_result = myrdd.map(f).collect()
    

    【讨论】:

      猜你喜欢
      • 1970-01-01
      • 1970-01-01
      • 1970-01-01
      • 2019-01-27
      • 2021-05-21
      • 1970-01-01
      • 1970-01-01
      • 1970-01-01
      • 1970-01-01
      相关资源
      最近更新 更多