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

如何按小时聚合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行。

当前处理代码如下:

  1. 将时间戳转换为小时级格式:
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|
  1. 创建行程计数列:
df_new_final=df_new_final.withColumn("Trip_count", lit(1))
  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|
  1. 聚合后行数:
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

相关产品推荐
方舟 Agent Plan

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

最近更新时间:2026.08.25 09:24:28