如何按小时聚合Spark DataFrame并展示含零计数的时段-地点组合
纽约黄色出租车数据集补全缺失时段行程计数问题
我正在处理纽约黄色出租车2020年数据集,包含近2400万条上下车记录,核心关注两列:tpep_pickup_datetime(上车时间戳)和PULocationID(上车地点ID),数据样例如下:
df_new_final.select(['tpep_pickup_datetime','PULocationID']).show() +--------------------+------------+ |tpep_pickup_datetime|PULocationID| +--------------------+------------+ | 2020-03-07 12:35:20| 79| | 2020-03-07 13:06:13| 107| | 2020-03-07 13:40:17| 264| | 2020-03-07 13:41:34| 236| | 2020-03-07 14:11:24| 230| | 2020-03-07 14:29:23| 239| | 2020-03-07 14:46:16| 230| | 2020-03-07 15:49:30| 170| | 2020-03-07 16:53:56| 261| | 2020-03-08 06:02:16| 114| | 2020-03-08 06:23:47| 142| | 2020-03-08 07:00:15| 186| | 2020-03-08 08:05:36| 170| | 2020-03-08 09:21:57| 148| | 2020-03-08 10:10:08| 13| | 2020-03-08 10:48:19| 162| | 2020-03-08 10:56:04| 233| | 2020-03-08 11:17:04| 170| | 2020-03-08 11:29:25| 162| | 2020-03-08 11:53:42| 138| +--------------------+------------+
问题描述
数据集共有262个地点ID,我的目标是按小时聚合每个地点的行程数量。2020年是闰年共366天,理论上所有时段-地点组合的总行数应为 366×24×262=2301408,但当前聚合后仅得到 910784行。
当前处理代码如下:
- 将时间戳转换为小时级格式:
df_new_final=df_new_final.withColumn("Pickup_datetime_hourly", date_format(col("tpep_pickup_datetime").cast("timestamp"), "yyyy-MM-dd HH:00"))
转换后数据样例:
+--------------------+----------------------+------------+ |tpep_pickup_datetime|Pickup_datetime_hourly|PULocationID| +--------------------+----------------------+------------+ | 2020-03-07 12:35:20| 2020-03-07 12:00| 79| | 2020-03-07 13:06:13| 2020-03-07 13:00| 107| | 2020-03-07 13:40:17| 2020-03-07 13:00| 264| | 2020-03-07 13:41:34| 2020-03-07 13:00| 236| | 2020-03-07 14:11:24| 2020-03-07 14:00| 230| | 2020-03-07 14:29:23| 2020-03-07 14:00| 239| | 2020-03-07 14:46:16| 2020-03-07 14:00| 230| | 2020-03-07 15:49:30| 2020-03-07 15:00| 170|
- 创建行程计数列:
df_new_final=df_new_final.withColumn("Trip_count", lit(1))
- 按小时和地点ID分组聚合:
hourly_aggregated=df_new_final.groupby(['Pickup_datetime_hourly','PULocationID']).agg({'Trip_count':'count'})
聚合后数据样例:
+----------------------+------------+-----------------+ |Pickup_datetime_hourly|PULocationID|count(Trip_count)| +----------------------+------------+-----------------+ | 2020-03-01 07:00| 230| 72| | 2020-03-01 10:00| 232| 5| | 2020-03-01 16:00| 100| 180| | 2020-03-01 16:00| 179| 4| | 2020-03-01 18:00| 129| 3| | 2020-03-03 05:00| 168| 2| | 2020-03-03 06:00| 186| 392| | 2020-03-04 01:00| 33| 1| | 2020-03-04 04:00| 112| 1| | 2020-03-04 05:00| 170| 55| | 2020-03-04 20:00| 211| 128| | 2020-03-04 20:00| 166| 91| | 2020-03-05 03:00| 132| 22|
- 聚合后行数:
hourly_aggregated.count() 910784
问题在于:无行程的时段-地点组合不会出现在聚合结果中,比如某地点的部分时段缺失:
+----------------------+------------+-----------------+ |Pickup_datetime_hourly|PULocationID|count(Trip_count)| +----------------------+------------+-----------------+ | 2020-03-01 00:00| 230| 72| | 2020-03-01 01:00| 230| 5| | 2020-03-01 03:00| 230| 180| | 2020-03-01 06:00| 230| 4| | 2020-03-01 07:00| 230| 3|
我需要补全这些缺失的时段,将计数显示为0,示例如下:
+----------------------+------------+-----------------+ |Pickup_datetime_hourly|PULocationID|count(Trip_count)| +----------------------+------------+-----------------+ | 2020-03-01 00:00| 230| 72| | 2020-03-01 01:00| 230| 5| | 2020-03-01 02:00| 230| 0| | 2020-03-01 03:00| 230| 4| | 2020-03-01 04:00| 230| 0| | 2020-03-03 05:00| 230| 0| | 2020-03-03 06:00| 230| 4| | 2020-03-01 07:00| 230| 3|
解决方案
核心思路是生成所有可能的小时-地点组合,再与聚合结果做左连接,将缺失的计数填充为0。具体步骤如下:
1. 生成2020年所有小时级时间序列
from pyspark.sql import functions as F # 生成2020年全年每小时的时间戳序列 all_hours = spark.sql(""" SELECT sequence(to_timestamp('2020-01-01 00:00:00'), to_timestamp('2020-12-31 23:00:00'), interval 1 hour) as hours """).select(F.explode("hours").alias("Pickup_datetime_hourly")) # 转换为与原数据一致的字符串格式 all_hours = all_hours.withColumn("Pickup_datetime_hourly", F.date_format("Pickup_datetime_hourly", "yyyy-MM-dd HH:00"))
2. 获取所有唯一的地点ID
# 从原数据中提取所有不重复的PULocationID all_locations = df_new_final.select("PULocationID").distinct()
3. 生成小时-地点的笛卡尔积(所有可能组合)
# 交叉连接得到所有小时和地点的组合 all_combinations = all_hours.crossJoin(all_locations)
4. 左连接聚合结果并填充0值
# 左连接保留所有组合,将缺失的计数替换为0 final_result = all_combinations.join( hourly_aggregated, on=["Pickup_datetime_hourly", "PULocationID"], how="left" ).withColumn( "count(Trip_count)", F.coalesce(F.col("count(Trip_count)"), F.lit(0)) ).orderBy("Pickup_datetime_hourly", "PULocationID")
5. 验证结果行数
final_result.count() # 结果应为2301408
这样处理后,所有时段-地点组合都会被保留,无行程的时段计数将显示为0,完全符合需求。
内容的提问来源于stack exchange,提问作者Emislim
相关产品推荐
相关产品推荐

