如何提取含连续10个及以上null的行及其前后行至新DataFrame
解决Spark DataFrame提取连续10+Null块及前后行的问题
核心思路
先通过窗口函数标记连续的Null块,筛选出长度≥10的块,再提取这些块的所有行,以及块的前一行(非Null行)和后一行(非Null行)。
步骤详解与代码实现
1. 标记连续Null块ID
用窗口函数为每一段连续的Null行分配唯一ID,非Null行ID设为null:
from pyspark.sql import Window import pyspark.sql.functions as F # 按行号排序的窗口 window_order_rn = Window.orderBy("rn") df = df.withColumn( "null_block_id", # 当当前行是Null且前一行非Null时,生成新的块ID F.when( (F.col("isSpnValue") == 1) & (F.lag("isSpnValue", 1, 0).over(window_order_rn) == 0), F.monotonically_increasing_id() ).otherwise( # 连续Null行继承前一行的块ID F.when(F.col("isSpnValue") == 1, F.lag("null_block_id").over(window_order_rn)) ) )
2. 筛选长度≥10的Null块
统计每个Null块的行数,过滤出符合条件的块:
# 统计每个Null块的行数 null_block_counts = df.filter(F.col("null_block_id").isNotNull()) \ .groupBy("null_block_id") \ .agg(F.count("*").alias("block_length")) \ .filter(F.col("block_length") >= 10) # 获取有效块的ID列表(用join替代isin,避免大数据量下的collect性能问题) valid_blocks = null_block_counts.select("null_block_id")
3. 提取目标行(Null块+前后行)
先获取有效块的行号范围,再生成需要保留的所有行号,最后关联原DataFrame:
# 获取每个有效块的最小/最大行号 block_ranges = df.join(valid_blocks, on="null_block_id", how="inner") \ .groupBy("null_block_id") \ .agg( F.min("rn").alias("min_rn"), F.max("rn").alias("max_rn") ) # 生成需要保留的行号:块内所有行 + 块前一行 + 块后一行 block_ranges = block_ranges.withColumn( "block_rns", F.sequence(F.col("min_rn"), F.col("max_rn")).cast("array<int>") ).withColumn( "prev_rn", F.col("min_rn") - 1 ).withColumn( "next_rn", F.col("max_rn") + 1 ) # 展开并去重目标行号 target_rns = block_ranges.select( F.explode(F.array_union(F.col("block_rns"), F.array(F.col("prev_rn"), F.col("next_rn")))).alias("rn") ).distinct() # 关联原DataFrame得到结果 result_df = df.join(target_rns, on="rn", how="inner")
现有方法的问题分析
- 代码1(lead/lag筛选):窗口范围未精准匹配连续Null块,容易误将短Null块的lead/lag值关联到长块计数,导致误捕。
- 代码2(groupBy spn_grp):
spn_grp仅在非Null行递增,Null行属于前一个非Null组,无法覆盖Null块的后一行(属于新的非Null组)。
内容的提问来源于stack exchange,提问作者IDK
相关产品推荐
相关产品推荐

