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

Spark DataFrame分组时补全缺失日期并填充Null/0值的最优实现方案

解决Spark DataFrame补全连续日期并填充缺失值的问题

问题场景

假设我们有原始的事件日志数据,通过聚合得到了按UserName和date分组的统计DataFrame,但存在缺失日期和部分用户在某日期无记录的情况,需要补全连续日期,并为缺失记录的数值列填充0,最终得到连续时间序列的完整统计结果。

原始聚合后的DataFrame:

+--------+----------+-----------+-------------------+-------------------+
|UserName|date      |NoLogPerDay|NoLogPer-1st-12-hrs|NoLogPer-2nd-12-hrs|
+--------+----------+-----------+-------------------+-------------------+
|B       |2021-08-11|2          |2                  |0                  |
|A       |2021-08-11|3          |2                  |1                  |
|B       |2021-08-13|1          |1                  |0                  |
+--------+----------+-----------+-------------------+-------------------+

期望的最终DataFrame:

+--------+----------+-----------+-------------------+-------------------+
|UserName|date      |NoLogPerDay|NoLogPer-1st-12-hrs|NoLogPer-2nd-12-hrs|
+--------+----------+-----------+-------------------+-------------------+
|B       |2021-08-11|2          |2                  |0                  |
|A       |2021-08-11|3          |2                  |1                  |
|B       |2021-08-12|0          |0                  |0                  |
|A       |2021-08-12|0          |0                  |0                  |
|B       |2021-08-13|1          |1                  |0                  |
|A       |2021-08-13|0          |0                  |0                  |
+--------+----------+-----------+-------------------+-------------------+

核心思路

要实现这个需求,关键是先构建所有用户和所有连续日期的笛卡尔积,再和原聚合结果做左连接,最后填充缺失值。不能在groupBy过程中直接实现,因为groupBy只会统计存在的分组,必须在聚合后补全缺失的分组。

完整解决方案代码

import datetime as dt
from pyspark.sql import functions as F
from pyspark.sql.types import StructType,StructField, StringType, IntegerType, TimestampType, DateType

# 原始数据
dict2 = [("2021-08-11 04:05:06", "A"),
        ("2021-08-11 04:15:06", "B"),
        ("2021-08-11 09:15:26", "A"),
        ("2021-08-11 11:04:06", "B"),
        ("2021-08-11 14:55:16", "A"),
        ("2021-08-13 04:12:11", "B"),
        ]
schema = StructType([
    StructField("timestamp", StringType(), True),
    StructField("UserName", StringType(), True),
])

# 创建原始DataFrame
sdf = spark.createDataFrame(data=dict2,schema=schema)

# 转换时间格式,提取date列
sdf1 = sdf.withColumn('timestamp', F.to_timestamp("timestamp", "yyyy-MM-dd HH:mm:ss")) \
    .withColumn('date', F.to_date("timestamp")) \
    .select('timestamp', 'date', 'UserName')

# 第一步:按用户和日期聚合统计
df_agg = sdf1.groupBy("UserName", "date").agg(
    F.sum(F.hour("timestamp").between(0, 23).cast("int")).alias("NoLogPerDay"),  # 修正:0-23覆盖全天
    F.sum(F.hour("timestamp").between(0, 11).cast("int")).alias("NoLogPer-1st-12-hrs"),
    F.sum(F.hour("timestamp").between(12, 23).cast("int")).alias("NoLogPer-2nd-12-hrs"),
).sort('date', 'UserName')

# 第二步:获取所有唯一用户和连续日期的笛卡尔积
# 获取所有唯一用户
all_users = df_agg.select("UserName").distinct()
# 获取日期范围
min_date = sdf1.select(F.min('date')).first()[0]
max_date = sdf1.select(F.max('date')).first()[0]
# 生成连续日期序列(Spark 2.4+支持sequence函数)
dates_df = spark.sql(f"""
    SELECT sequence(to_date('{min_date}'), to_date('{max_date}'), interval 1 day) as dates
""").select(F.explode("dates").alias("date"))
# 生成用户-日期的全量组合
full_user_dates = all_users.crossJoin(dates_df)

# 第三步:左连接聚合结果,填充缺失值
final_df = full_user_dates.join(df_agg, on=["UserName", "date"], how="left") \
    .na.fill(0, subset=["NoLogPerDay", "NoLogPer-1st-12-hrs", "NoLogPer-2nd-12-hrs"]) \
    .sort("date", "UserName")

final_df.show(truncate=False)

关键步骤说明

  • 聚合统计:先完成基础的按用户和日期的统计,这里注意把hour between(0,24)修正为0-23,避免逻辑错误。
  • 生成全量分组:
    • 通过distinct()获取所有唯一用户
    • 使用Spark的sequence函数生成连续日期序列(比Python列表更高效,避免数据量过大时的性能问题)
    • 用crossJoin生成所有用户和所有日期的笛卡尔积,确保没有遗漏的分组
  • 左连接+填充:将全量分组和聚合结果左连接,然后用na.fill(0)为数值列的缺失值填充0,得到完整的结果。

为什么不建议在groupBy前处理?

如果在groupBy前补全日期,需要为每个缺失日期的每个用户添加空的事件记录,这会大幅增加数据量,尤其是用户和日期范围较大时,性能会很差。而在聚合后补全分组的方式更高效,只处理统计后的小数据量。

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

相关产品推荐
方舟 Agent Plan

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

最近更新时间:2026.05.01 00:03:11