【发布时间】:2018-04-05 03:18:39
【问题描述】:
我有一个带有如下列的 spark 数据框:
df
--------------------------
A B C D E F amt
"A1" "B1" "C1" "D1" "E1" "F1" 1
"A2" "B2" "C2" "D2" "E2" "F2" 2
我想使用列组合执行 groupBy
(A, B, sum(amt))
(A, C, sum(amt))
(A, D, sum(amt))
(A, E, sum(amt))
(A, F, sum(amt))
使得生成的数据框看起来像:
df_grouped
----------------------
A field value amt
"A1" "B" "B1" 1
"A2" "B" "B2" 2
"A1" "C" "C1" 1
"A2" "C" "C2" 2
"A1" "D" "D1" 1
"A2" "D" "D2" 2
我的尝试如下:
val cols = Vector("B","C","D","E","F")
//code for creating empty data frame with structs for the cols A, field, value and act
for (col <- cols){
empty_df = empty_df.union (df.groupBy($"A",col)
.agg(sum(amt).as(amt)
.withColumn("field",lit(col)
.withColumnRenamed(col, "value"))
}
我觉得“for”或“foreach”的用法对于像 spark 这样的分布式环境可能很笨拙。对于我正在做的事情,地图功能是否有任何替代方案?在我看来, aggregateByKey 和 collect_list 可能有效;但是,我无法想象一个完整的解决方案。请指教。
【问题讨论】:
-
您只是想取消旋转
B,C,D,E,F的值吗?sum(amt)在扮演什么角色? -
我的原始数据框是一个大集合。为了使其简单易懂,我将其压缩为几行。由于它很大,我认为 for 循环在内存使用方面可能不是最好的方法,所以 foldleft 可能会更好。
标签: scala apache-spark group-by