You need to enable JavaScript to run this app.
优惠活动
大模型
产品
解决方案
定价
更多

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()

小时级输出示例:

CAPTUREDTIMENODECHANNELLOCATIONTACKCOUNT
20-05-09 03:00:00PUSC_RESSIMPLEXNORTH_ALUE2200341
20-05-09 04:00:00PUSC_RESSIMPLEXSOUTH_ALUE2200342
20-05-09 12:00:00TESC_RESSIMPLEXNORTH_ALUE2200571

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()

日级输出示例:

CAPTUREDTIMENODECHANNELLOCATIONTACKCOUNT
20-05-09 00:00:00PUSC_RESSIMPLEXNORTH_ALUE2200341
20-05-09 00:00:00PUSC_RESSIMPLEXSOUTH_ALUE2200341
20-05-09 00:00:00TESC_RESSIMPLEXNORTH_ALUE2200571

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()

周级输出示例:

CAPTUREDTIMENODECHANNELLOCATIONTACKCOUNT
20-05-03 00:00:00PUSC_RESSIMPLEXNORTH_ALUE2200341
20-05-03 00:00:00TESC_RESSIMPLEXNORTH_ALUE2200573

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()

月级输出示例:

CAPTUREDTIMENODECHANNELLOCATIONTACKCOUNT
20-04-01 00:00:00TESC_RESSIMPLEXNORTH_ALUE2200572
20-04-01 00:00:00PUSC_RESSIMPLEXNORTH_ALUE2200711
20-05-01 00:00:00PUSC_RESSIMPLEXNORTH_ALUE2200341

关键说明

  • date_trunc是PySpark中做时间粒度截断的核心函数,支持hour/day/week/month等多种粒度参数。
  • 如果你的CAPTUREDTIME已经是Timestamp类型,可以跳过第一步的转换操作。
  • date_format用来把截断后的时间戳转成你需要的字符串格式,确保输出和示例一致。

内容的提问来源于stack exchange,提问作者stacktesting

相关产品推荐
方舟 Agent Plan

超全模态模型 × Harness 升级,最新支持 Deepseek-V4.1-Flash、GLM-5.3 系列、Doubao-Seedream-5.0-pro、Kimi-K3 (部分), 限时 9.9 元起

最近更新时间:2026.05.08 17:32:30