分组操作中如何获取DataFrame的最早/最晚日期及对应用户等统计信息
问题分析与解决方案
你当前的代码存在几个关键问题,导致无法正确获取每个url下最早date对应的user:
- 窗口函数的分区逻辑错误:你按
url和user分区,这会给每个用户在对应url下单独排序,无法找到整个url维度下最早日期的用户。 - 语法错误:
agg方法里的sum("followers")多了一个多余的引号,会导致编译失败。 - 操作顺序错误:
groupBy之后的DataFrame已经丢失了user和date字段,后续无法基于这些字段使用窗口函数。
下面是两种可行的解决方案,你可以根据需求选择:
方案一:先提取最早用户再聚合(清晰易懂)
这种方式先单独获取每个url的最早用户和日期,再和聚合结果关联,逻辑更直观:
import org.apache.spark.sql.functions.{sum, countDistinct, min, max, row_number} import org.apache.spark.sql.Window // 1. 创建窗口:按url分区,按date升序排序,找到每个url下最早的行 val urlWindow = Window.partitionBy("url").orderBy($"date".asc) // 2. 过滤出每个url的第一行,得到最早用户和日期 val firstUserDF = df.withColumn("rn", row_number().over(urlWindow)) .filter($"rn" === 1) .select("url", "user", "date") .withColumnRenamed("user", "first_user") .withColumnRenamed("date", "earliest_date") // 3. 分组聚合其他统计指标 val aggDF = df.groupBy("url") .agg( sum("followers").alias("sum_followers"), countDistinct("user").alias("distinct_users"), max("date").alias("latest_date") ) // 4. 关联两个结果,得到最终数据 val result = aggDF.join(firstUserDF, Seq("url"), "inner") .select("url", "first_user", "earliest_date", "latest_date", "sum_followers", "distinct_users")
方案二:单窗口+聚合(更高效)
这种方式先给每个url下的所有行打上最早用户和日期的标记,再直接分组聚合,避免了join操作:
import org.apache.spark.sql.functions.{sum, countDistinct, max, first} import org.apache.spark.sql.Window // 创建窗口:按url分区,按date升序排序 val urlWindow = Window.partitionBy("url").orderBy($"date".asc) // 给每行添加最早用户和日期(同一url下所有行的这两个值一致) val enrichedDF = df.withColumn("first_user", first($"user").over(urlWindow)) .withColumn("earliest_date", first($"date").over(urlWindow)) // 直接分组聚合所有需要的指标 val result = enrichedDF.groupBy("url") .agg( first("first_user").alias("first_user"), first("earliest_date").alias("earliest_date"), max("date").alias("latest_date"), sum("followers").alias("sum_followers"), countDistinct("user").alias("distinct_users") )
验证结果
运行上述代码后,你会得到符合预期的输出(注意你给出的预期输出中www.test2.com的日期顺序有误,正确的最早日期是2017-01-03 17:00:00,最晚是2017-01-04 15:05:00):
| url | first_user | earliest_date | latest_date | sum_followers | distinct_users |
|---|---|---|---|---|---|
| www.test1.com | B | 2017-01-03 10:46:00 | 2017-01-05 05:46:00 | 66 | 3 |
| www.test2.com | B | 2017-01-03 17:00:00 | 2017-01-04 15:05:00 | 55 | 2 |
内容的提问来源于stack exchange,提问作者ScalaBoy
相关产品推荐
相关产品推荐

