【问题标题】:Running a function in the last RDD in spark在 spark 的最后一个 RDD 中运行一个函数
【发布时间】:2020-05-05 07:18:34
【问题描述】:

我有一个 python 程序来分析数据并想用 Spark 运行它。我在工人之间分配数据并对其进行一些转换。但最后我需要将结果收集到主节点并在其上运行另一个函数。

在驱动程序上我有这个代码:

sc = SparkContext(conf=spark_conf)
sc.parallelize(group_list, 4) \
    .map(function1, preservesPartitioning=True) \
    .map(function2, preservesPartitioning=True) \
    .map(function3, preservesPartitioning=True) \
    .map(function4, preservesPartitioning=True) \
    .map(function5, preservesPartitioning=True) \
    .map(function6, preservesPartitioning=True) \
    .map(function7, preservesPartitioning=True) \
    .map(function8, preservesPartitioning=True) \
    .map(function9, preservesPartitioning=True)

function9 生成的最后一个 RDD 是一个包含多行和唯一键的表。当主节点从worker收集所有最后一个RDD时,它们在主节点中有重复的行。我必须按最后一个表进行分组并对某些列进行一些聚合,所以我有一个最终函数,它采用最后一个表并对其进行分组和聚合。但是我不知道如何将最后一个 RDD 传递给 final 函数。

例如在worker1上,我有这个数据:

    key    count   average
     B       3       0.2
     x       2       0.1
     y       5       1.2

在worker2上,我有这个数据:

    key    count    average
     B       2         0.1
     c       1         0.01
     x       3         0.34

当主节点从worker接收到所有数据时,它有:

    key    count    average
     B       3       0.2
     x       2       0.1
     y       5       1.2
     B       2       0.1
     c       1       0.01
     x       3       0.34

你看到数据有两个 B 和两个 x 键。我必须在主节点中使用另一个函数按 key 列进行分组并计算 average 列的新平均值。我使用了 reduce 并将我的最终函数赋予它,但它给了我错误,因为它需要两个参数。 请指导我可以使用什么 spark 操作在最后一个 RDD 上运行我的函数?

非常感谢任何指导。

【问题讨论】:

    标签: python apache-spark pyspark


    【解决方案1】:

    我建议你传递给 DataFrame 格式(使用起来更简单),然后应用这个:

    df.groupBy('key').agg(f.sum('count'), f.avg('average'))
    

    如果您想保持 rdd 格式,您应该执行类似 this 的操作,但应用平均值而不是列表。

    从你写的应该可以:

    sqlContext = sql.SQLContext(sc)
    from pyspark.sql import SQLContext
    
     (sqlContext.createDataFrame(
         [['B',3,0.2],
          ['x',2,0.1],
          ['y',5,1.2],
          ['B',2,0.1],
          ['c',1,0.01],
          ['x',3,0.34]], ['key', 'count', 'average'])
     .groupBy('key')
     .agg(f.sum('count').alias('count'), f.avg('average').alias('avg'))
     .show()
    )
    

    您可以(并且可能应该)传递初始 rdd sc.parallelize(group_list, 4),在这种情况下,f.sum() 应该是 f.count()。希望这会有所帮助

    【讨论】:

    • 亲爱的@ggagliano,感谢您的反馈,但我知道 df.goupby。问题是如何让master中的所有worker rdd在master节点中运行df.groupby
    • 如果我理解得很好,我认为这不是使用 spark 的最有效方式。您应该只将聚合的最终结果检索到主服务器。我要做的是 groupBy..agg 然后收集到主节点。但是,如果您确信您的行很少,您可以在最终聚合之前收集,比如说,转换为 pandas 数据帧并在单个节点上进行聚合。希望这会有所帮助
    • 你的意思是,如果我在程序中执行 collect() 操作并将结果转换为 pandas Dataframe,我可以执行 groupby 在主节点上?非常感谢您的指导。
    • 您可以直接使用 df.toPandas() 收集为 Pandas Dataframe。它不会完全在 spark 主节点上运行,因为您要在 spark 之外运行。它将在您当前的 python 解释器上运行,因此,如果您对行的基数没有信心,我不建议您这样做。在这种情况下,最好执行 df.limit(10000).toPandas() 之类的操作,以确保不收集十亿行
    • 事实上,我之前对每个工人都做了 df.groupby,但由于我无法将所有十亿行发送给工人,我不得不将它们分配给工人。因此,当我从工人那里收集所有数据时,我在主节点中有重复的行,因此我必须再次在主节点中执行 df.groupby 以删除重复的行。请您在新帖子中写下您的答案以选择它作为接受的答案吗?
    【解决方案2】:

    例如,我有一个这样的熊猫数据框:

    'a'     'b'
     1       3
     1       4
     2       5
    

    我写了一个在 groupby 中使用的函数:

    def process_json(x):
    print(x)
    temp = 0
    for item in x.items():
        temp += item[1]
    print('-------', temp)
    

    因此,

    a.groupby(['a'])['b'].agg(process_json)
    

    输出是:

    0    3
    1    4
    Name: b, dtype: int64
    ------- 7
    2    5
    Name: b, dtype: int64
    ------- 5
    

    【讨论】:

      猜你喜欢
      • 1970-01-01
      • 2016-12-23
      • 2021-10-25
      • 1970-01-01
      • 2017-08-19
      • 1970-01-01
      • 1970-01-01
      • 2016-03-21
      • 2020-02-12
      相关资源
      最近更新 更多