如何用Spark操作MongoDB聚合 将特殊用户“-”的data1值叠加到其他用户上
实现思路
整体逻辑分3步即可完成需求:
- 从分组后的结果集中,单独提取出用户名为
-的特殊记录对应的data1值 - 遍历所有记录,对非
-的普通用户,将其data1值与提取出的特殊值叠加,特殊记录本身保持原值不变 - 输出处理后的结果即可
代码示例
Scala版本
假设你groupBy后聚合得到的DataFrame结构为(username: String, data1: Int):
import org.apache.spark.sql.functions._ // groupedDF为你之前分组聚合完成的数据集 val groupedDF = rawDataRows.groupBy("username").agg(sum("data1").alias("data1")) // 提取特殊值,数据量较大时可以广播该常量优化性能 val specialVal = groupedDF.filter(col("username") === "-").select("data1").head().getAs[Int](0) val broadcastVal = spark.sparkContext.broadcast(specialVal) // 叠加计算 val resultDF = groupedDF.withColumn("data1", when(col("username") =!= "-", col("data1") + broadcastVal.value) .otherwise(col("data1")) ) // 查看结果 resultDF.show()
PySpark版本
from pyspark.sql.functions import when # grouped_df为你之前分组聚合完成的数据集 grouped_df = rawDataRows.groupBy("username").sum("data1").withColumnRenamed("sum(data1)", "data1") # 提取特殊值 special_val = grouped_df.filter(grouped_df.username == "-").select("data1").head()[0] # 叠加计算 result_df = grouped_df.withColumn("data1", when(grouped_df.username != "-", grouped_df.data1 + special_val) .otherwise(grouped_df.data1) ) # 查看结果 result_df.show()
注意事项
- 若业务场景中可能不存在用户名为
-的记录,需要在提取特殊值时增加判空逻辑,避免空指针异常 - 若使用RDD实现,逻辑完全一致:先过滤出特殊记录取值,再map所有元素叠加即可
内容的提问来源于stack exchange,提问作者RemovedQuasar
相关产品推荐
相关产品推荐

