Spark中如何将时间戳数据聚合为小时粒度 按位置统计上车量
Spark实现纽约出租车数据按位置小时粒度聚合方案
你用到的纽约出租车数据集字段包含上车位置PULocationID、上车时间戳tpep_pickup_datetime,数据样例如下:
纽约出租车数据集样例
核心实现逻辑:将时间戳字段向下截断到小时维度,与上车位置共同作为分组键,统计每组的记录总数即为对应位置每小时的上车量,以下是两种生产环境常用的实现方式:
方式1:PySpark DataFrame API 实现(推荐)
- 首先加载源数据,确保上车时间字段为
Timestamp类型,若为字符串格式需先做类型转换 - 调用
date_trunc内置函数将时间截断到小时粒度,生成小时维度列 - 按上车位置ID、小时维度列分组,聚合统计组内记录数
- 可根据业务需要对结果排序、写入持久化存储
对应实现代码:
from pyspark.sql import SparkSession from pyspark.sql.functions import date_trunc, count, to_timestamp # 初始化Spark会话 spark = SparkSession.builder.appName("nyc_taxi_hourly_pickup_stat").getOrCreate() # 替换为你实际的数据源路径,支持parquet、csv等Spark兼容格式 source_df = spark.read.parquet("/path/to/your/nyc_taxi_source_data") # 若tpep_pickup_datetime为字符串类型,取消注释下一行执行类型转换 # source_df = source_df.withColumn("tpep_pickup_datetime", to_timestamp("tpep_pickup_datetime")) # 核心聚合逻辑 hourly_pickup_stat = source_df.withColumn( "pickup_hour", date_trunc("hour", "tpep_pickup_datetime") ).groupBy( "PULocationID", "pickup_hour" ).agg( count("*").alias("pickup_total_cnt") ) # 预览前20条结果 hourly_pickup_stat.show() # 可选:将结果写入存储 # hourly_pickup_stat.write.mode("overwrite").parquet("/path/to/output/hourly_pickup_res")
注意:如果数据集时间戳为UTC时区,需要匹配本地业务时区的话,先使用
from_utc_timestamp函数转换时区后再做时间截断,避免小时维度统计错位。
方式2:Spark SQL 实现
如果习惯使用SQL语法,可将DataFrame注册为临时视图后直接执行SQL查询,逻辑与API完全一致:
# 注册临时视图 source_df.createOrReplaceTempView("nyc_taxi_trip") hourly_pickup_stat_sql = spark.sql(""" SELECT PULocationID, date_trunc('hour', tpep_pickup_datetime) AS pickup_hour, COUNT(*) AS pickup_total_cnt FROM nyc_taxi_trip GROUP BY PULocationID, date_trunc('hour', tpep_pickup_datetime) """) hourly_pickup_stat_sql.show()
补充说明:如果后续需要调整聚合粒度,只需要修改date_trunc函数的第一个粒度参数即可,支持minute、day、week、month等常用时间粒度,无需改动其他聚合逻辑。
内容的提问来源于stack exchange,提问作者Emislim
相关产品推荐
相关产品推荐

