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

如何基于起止日期合并PySpark DataFrame同组连续行?

处理十亿级PySpark DataFrame的连续日期区间合并

针对十亿条记录的大规模数据,要高效实现按key1/key2分组、合并连续日期区间的需求,可以通过窗口函数+分组聚合的方式完成,具体步骤如下:

1. 基础排序与窗口定义

首先按key1、key2对数据分区,组内按effective_start升序排序,这是后续对比相邻行日期的基础:

from pyspark.sql import Window
import pyspark.sql.functions as F

# 定义窗口规则:按key1/key2分区,按effective_start排序
window_spec = Window.partitionBy("key1", "key2").orderBy("effective_start")

2. 标记连续区间分界

用lag函数获取同组内前一行的effective_end,对比当前行的effective_start,生成标记列区分是否属于同一连续区间:

  • 若当前行effective_start等于前一行effective_end,标记为0(同区间)
  • 否则标记为1(新区间起点)
df_with_lag = df.withColumn(
    "prev_end",
    F.lag("effective_end").over(window_spec)
).withColumn(
    "is_continuous",
    F.when(F.col("effective_start") == F.col("prev_end"), 0).otherwise(1)
)

3. 生成连续区间组ID

对is_continuous列做累计求和,同一连续区间内的行将获得相同的group_id,不同区间的ID自动递增:

df_with_group = df_with_lag.withColumn(
    "group_id",
    F.sum("is_continuous").over(window_spec.rangeBetween(Window.unboundedPreceding, 0))
)

4. 聚合合并区间

按key1、key2、group_id分组,取每组内最小的effective_start和最大的effective_end,得到最终合并结果:

result_df = df_with_group.groupBy("key1", "key2", "group_id").agg(
    F.min("effective_start").alias("merged_start"),
    F.max("effective_end").alias("merged_end")
).drop("group_id")

十亿级数据性能优化建议

  • 预分区优化:执行df.repartition("key1", "key2")让相同key组合的数据落在同一分区,减少窗口操作的shuffle开销。
  • 分桶表持久化:若数据长期复用,可预先创建按key1、key2分桶的表,大幅降低后续分组、窗口操作的资源消耗。
  • 精简列数据:处理过程中仅保留key1、key2、effective_start、effective_end必要列,减少内存占用与数据传输量。

针对你提到的key11/key2、key22/key2、key22/key3三个分组,执行完上述步骤后,每组内的连续日期行都会被合并为对应的区间行。

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

相关产品推荐
方舟 Agent Plan

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

最近更新时间:2026.08.11 05:50:25