PySpark DataFrame按时间频率分组统计问题
解决PySpark按多时间粒度分组聚合的问题
嘿,刚上手PySpark的话,按不同时间粒度分组确实得先把时间字段的类型搞定——你现在的CAPTUREDTIME看起来是字符串格式,第一步得把它转成PySpark能识别的时间戳类型,不然没法做时间截断的操作。我给你一步步拆解怎么实现小时、日、周、月级的分组聚合:
第一步:转换时间字段类型
先把字符串格式的CAPTUREDTIME转成PySpark的TimestampType,这样才能用时间相关的函数做分组:
from pyspark.sql import functions as func from pyspark.sql.types import TimestampType # 转换时间格式(你的时间格式是yy-MM-dd HH:mm:ss,要对应指定) df = df.withColumn("CAPTUREDTIME", func.to_timestamp("CAPTUREDTIME", "yy-MM-dd HH:mm:ss"))
1. 小时级(Hourly)分组
用date_trunc("hour", ...)把时间截断到当前小时的起始时刻(比如03:06:21变成03:00:00),然后和其他维度字段一起分组,最后聚合计数:
# 小时级分组 hourly_df = df.groupBy( func.date_trunc("hour", "CAPTUREDTIME").alias("CAPTUREDTIME"), "NODE", "CHANNEL", "LOCATION", "TACK" ).agg(func.count("TACK").alias("COUNT")) # 把时间戳转成你需要的字符串格式 hourly_df = hourly_df.withColumn("CAPTUREDTIME", func.date_format("CAPTUREDTIME", "yy-MM-dd HH:mm:ss")) hourly_df.show()
小时级输出示例:
| CAPTUREDTIME | NODE | CHANNEL | LOCATION | TACK | COUNT |
|---|---|---|---|---|---|
| 20-05-09 03:00:00 | PUSC_RES | SIMPLEX | NORTH_AL | UE220034 | 1 |
| 20-05-09 04:00:00 | PUSC_RES | SIMPLEX | SOUTH_AL | UE220034 | 2 |
| 20-05-09 12:00:00 | TESC_RES | SIMPLEX | NORTH_AL | UE220057 | 1 |
2. 日级(Daily)分组
同理,用date_trunc("day", ...)把时间截断到当日00:00:00:
# 日级分组 daily_df = df.groupBy( func.date_trunc("day", "CAPTUREDTIME").alias("CAPTUREDTIME"), "NODE", "CHANNEL", "LOCATION", "TACK" ).agg(func.count("TACK").alias("COUNT")) daily_df = daily_df.withColumn("CAPTUREDTIME", func.date_format("CAPTUREDTIME", "yy-MM-dd HH:mm:ss")) daily_df.show()
日级输出示例:
| CAPTUREDTIME | NODE | CHANNEL | LOCATION | TACK | COUNT |
|---|---|---|---|---|---|
| 20-05-09 00:00:00 | PUSC_RES | SIMPLEX | NORTH_AL | UE220034 | 1 |
| 20-05-09 00:00:00 | PUSC_RES | SIMPLEX | SOUTH_AL | UE220034 | 1 |
| 20-05-09 00:00:00 | TESC_RES | SIMPLEX | NORTH_AL | UE220057 | 1 |
3. 周级(Weekly)分组
用date_trunc("week", ...)把时间截断到当周的起始时刻(默认是周日,如果你需要周一作为一周的开始,可以用func.date_trunc("week", func.date_add("CAPTUREDTIME", 1))调整):
# 周级分组(默认周日为一周起始) weekly_df = df.groupBy( func.date_trunc("week", "CAPTUREDTIME").alias("CAPTUREDTIME"), "NODE", "CHANNEL", "LOCATION", "TACK" ).agg(func.count("TACK").alias("COUNT")) weekly_df = weekly_df.withColumn("CAPTUREDTIME", func.date_format("CAPTUREDTIME", "yy-MM-dd HH:mm:ss")) weekly_df.show()
周级输出示例:
| CAPTUREDTIME | NODE | CHANNEL | LOCATION | TACK | COUNT |
|---|---|---|---|---|---|
| 20-05-03 00:00:00 | PUSC_RES | SIMPLEX | NORTH_AL | UE220034 | 1 |
| 20-05-03 00:00:00 | TESC_RES | SIMPLEX | NORTH_AL | UE220057 | 3 |
4. 月级(Monthly)分组
用date_trunc("month", ...)把时间截断到当月1号的00:00:00:
# 月级分组 monthly_df = df.groupBy( func.date_trunc("month", "CAPTUREDTIME").alias("CAPTUREDTIME"), "NODE", "CHANNEL", "LOCATION", "TACK" ).agg(func.count("TACK").alias("COUNT")) monthly_df = monthly_df.withColumn("CAPTUREDTIME", func.date_format("CAPTUREDTIME", "yy-MM-dd HH:mm:ss")) monthly_df.show()
月级输出示例:
| CAPTUREDTIME | NODE | CHANNEL | LOCATION | TACK | COUNT |
|---|---|---|---|---|---|
| 20-04-01 00:00:00 | TESC_RES | SIMPLEX | NORTH_AL | UE220057 | 2 |
| 20-04-01 00:00:00 | PUSC_RES | SIMPLEX | NORTH_AL | UE220071 | 1 |
| 20-05-01 00:00:00 | PUSC_RES | SIMPLEX | NORTH_AL | UE220034 | 1 |
关键说明
date_trunc是PySpark中做时间粒度截断的核心函数,支持hour/day/week/month等多种粒度参数。- 如果你的
CAPTUREDTIME已经是Timestamp类型,可以跳过第一步的转换操作。 date_format用来把截断后的时间戳转成你需要的字符串格式,确保输出和示例一致。
内容的提问来源于stack exchange,提问作者stacktesting
相关产品推荐
相关产品推荐

