如何在PySpark DataFrame中实现干小时数递增与降雨时归零
解决PySpark计算连续干小时数的问题
你的需求是计算连续无降雨小时数:无降雨时数值递增,有降雨时重置为0。原代码的问题在于未正确对连续干小时区间分段,全局累加会导致跨降雨区间的错误计数,同时逻辑运算符使用有误。
正确实现思路
核心是先将连续干小时划分为独立分组,再在每个分组内计算连续计数:
- 标记每小时是否降雨(1=有雨,0=无雨)
- 通过累加历史降雨标记生成分组ID,每次遇到降雨时分组ID递增,实现连续干小时的分段
- 在每个分组内按时间排序,计算行号,无降雨时行号即为干小时数,有降雨时置0
完整代码实现
from pyspark.sql import functions as F from pyspark.sql.window import Window def calculate_dry_hours(df): # 1. 标记是否降雨:1=有雨,0=无雨 df = df.withColumn("rain_flag", F.when(F.col("rain_1h") > 0, 1).otherwise(0)) # 2. 生成分组键:累加历史降雨标记,实现连续干小时的分段 window_order = Window.orderBy("dt") df = df.withColumn("group_id", F.sum("rain_flag").over(window_order)) # 3. 在每个分组内计算连续干小时数 window_group = Window.partitionBy("group_id").orderBy("dt") df = df.withColumn( "dry_hours", F.when( F.col("rain_flag") == 0, F.row_number().over(window_group) ).otherwise(0) ) # 可选:删除中间辅助列 df = df.drop("rain_flag", "group_id") return df
代码解释
- rain_flag:明确标记当前小时是否有降雨,为后续分组做准备
- group_id:通过累加所有历史的
rain_flag值,每次遇到降雨(rain_flag=1)时,group_id会增加1,这样连续的干小时会被分到同一个group_id下,降雨小时则会开启新分组 - dry_hours:在每个
group_id分组内,用row_number()计算当前是该分组的第几个干小时,无降雨时返回行号,降雨时直接置0
示例验证
假设输入数据:
| dt | rain_1h |
|---|---|
| 2024-01-01 00:00:00 | 0 |
| 2024-01-01 01:00:00 | 0 |
| 2024-01-01 02:00:00 | 5.2 |
| 2024-01-01 03:00:00 | 0 |
| 2024-01-01 04:00:00 | 0 |
| 2024-01-01 05:00:00 | 0 |
运行代码后输出:
| dt | rain_1h | dry_hours |
|---|---|---|
| 2024-01-01 00:00:00 | 0 | 1 |
| 2024-01-01 01:00:00 | 0 | 2 |
| 2024-01-01 02:00:00 | 5.2 | 0 |
| 2024-01-01 03:00:00 | 0 | 1 |
| 2024-01-01 04:00:00 | 0 | 2 |
| 2024-01-01 05:00:00 | 0 | 3 |
内容的提问来源于stack exchange,提问作者Luis Valencia
相关产品推荐
相关产品推荐

