【发布时间】: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