我有兩列的火花資料幀
colA colB
1 3
1 2
2 4
2 5
2 1
我想 groupBy colA 并迭代每個組的 colB 串列,以便:
res = 0
for i in collect_list(col("colB")):
res = i * (3 res)
回傳值應為 res
所以我得到:
colA colB
1 24
2 78
我怎樣才能在 Scala 中做到這一點?
uj5u.com熱心網友回復:
您可以通過以下方式獲得您想要的結果:
val df = Seq((1,3), (1,2), (2,4), (2,5), (2,1)).toDF("colA", "colB")
val retDf = df
.groupBy("colA")
.agg(
aggregate(
collect_list("colB"), lit(0), (acc, nxt) => nxt * (acc 3)
) as "colB")
但是,您需要非常小心,因為 Spark 上的資料是分布式的。如果資料在被讀入 Spark 后被打亂,則不能保證它會以相同的順序被收集。在玩具示例collect_list("colB")中將回傳Seq(3,2)where colAis 1。但是,如果在較早的階段進行了任何洗牌,collect_list也可以回傳Seq(2,3),這將給您27而不是所需的24. 您需要為您的資料提供一些元資料,您可以使用這些元資料來確保按照您期望的順序處理這些資料,例如使用monotonicallyIncreasingId方法。
uj5u.com熱心網友回復:
RDD 方法,不會丟失排序。
%scala
val rdd1=spark.sparkContext.parallelize(Seq((1,3), (1,3), (2,4), (2,5), (2,1))).zipWithIndex().map(x => ((x._1._1), (x._1._2, x._2)) )
val rdd2 = rdd1.groupByKey
// Convert to Array.
val rdd3 = rdd2.map(x => (x._1, x._2.toArray))
val rdd4 = rdd3.map(x => (x._1, x._2.sortBy(_._2)))
val rdd5 = rdd4.mapValues(v => v.map(_._1))
rdd5.collect()
val res = rdd5.map(x => (x._1, x._2.fold(0)((acc, nxt) => nxt * (acc 3) )))
res.collect()
回傳:
res201: Array[(Int, Int)] = Array((1,24), (2,78))
根據需要從 DF 到 DF 進行隱蔽。
轉載請註明出處,本文鏈接:https://www.uj5u.com/houduan/322630.html
上一篇:Spark結構化流批量讀取檢查點
