【发布时间】:2016-12-07 23:52:21
【问题描述】:
如果我们使用.reduce(max),那么我们将得到整个RDD中最大的key。我知道这个 reduce 将在所有分区上运行,然后减少每个分区发送的那些项目。但是我们如何才能取回每个分区的最大键呢?为.mapPartitions()写一个函数?
【问题讨论】:
标签: apache-spark pyspark apache-spark-sql spark-streaming
如果我们使用.reduce(max),那么我们将得到整个RDD中最大的key。我知道这个 reduce 将在所有分区上运行,然后减少每个分区发送的那些项目。但是我们如何才能取回每个分区的最大键呢?为.mapPartitions()写一个函数?
【问题讨论】:
标签: apache-spark pyspark apache-spark-sql spark-streaming
你可以:
rdd.mapParitions(iter => Iterator(iter.reduce(Math.max)))
或
rdd.mapPartitions(lambda iter: [max(iter)])
在流式传输中使用 DStream.trasform。
【讨论】: