【发布时间】:2021-02-17 08:35:30
【问题描述】:
我有一个像这样的 PySpark 数据框:
name | A | Number
--------------------
Dan | 1 | 1.5
Jan | 1 | 2.5
Dan | 1 | 1.5
Ann | 1 | 1.5
Ann | 2 | 1.2
Jon | 2 | 1.7
在数字列中有一个名称列和一个与之关联的值。名称将始终有一个对应的数字,即如果 Dan 的值为 1.5,那么无论 Dan 出现在数据中的何处,关联的数字都是 1.5。
我想在 A 列上进行分组,并找到唯一名称的 Number 的不同总和。
例如,对于第 1 组,唯一名称是 Dan、Jan 和 Ann。它们对应的数字是 1.5、2.5、1.5,总和为 5.5。
我试图通过创建一个虚拟列来解决这个问题,将其按组收集到一组,然后将其分解并求和。代码是这样的,
sdf = sdf.withColumn("dummy_col", F.concat(F.col("name"), F.lit("_"), F.col("Number")))
sdf = sdf.groupby("A") \
.agg(F.collect_set("dummy_col").alias("dummy_set"))
sdf = sdf.withColumn("exploded", F.explode("dummy_set")) \
.withColumn("Number_new", F.split("exploded")[1])
sdf_final = sdf.groupby("A") \
.agg(F.sum("Number_new").alias("Sum"))
但这个解决方案根本不可扩展。当数据为数百万行时,运行时间会很长。
我正在运行 spark 2.3,我也无法升级它。
我正在寻找一种能够很好地适应数据大小并且步骤更少的解决方案。另外,我必须在 A、B、C 和 D 等几列上做一个多维数据集,而不是在 A 上进行 groupby 并做类似的求和。
我想要的是运行一个多维数据集而不是 groupby。 drop_duplicates 在那里帮不了我。另外,我的键是很多列,而不仅仅是 A 列。
作为 cube 的替代品,我可以创建一个 groupby 级别列表并运行一个循环并一个一个地处理它以使用 drop_duplicates。
这仍然会导致与性能相关的问题。我还没有尝试过。但我可以看到 groupby 键的列表至少有 100 项。
【问题讨论】:
标签: python python-3.x pyspark