如何基于起止日期合并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
相关产品推荐
相关产品推荐

