窗口嵌套创建方法及DataFrame分区窗口数据填充技术咨询
嘿,我来帮你搞定这两个问题!
问题1:如何在一个窗口(Window)内创建另一个窗口?
其实不管是用PySpark还是Pandas,核心思路都是先定义一个大的主窗口(比如按name分区),然后在这个主窗口内部,通过标记分组边界的方式,拆分出更小的子窗口。
拿PySpark举例子:
首先你得先定义好主窗口,比如按name分区并按顺序排列:
from pyspark.sql import Window import pyspark.sql.functions as F main_window = Window.partitionBy("name").orderBy("row_num") # row_num是用来保证顺序的列
接下来,你需要找到子窗口的边界——比如你第二个问题里的“非1值”就是边界点。我们可以先给这些边界点打标记,然后通过累加标记来生成每个子窗口的ID:
# 给非1值的行打个标记 df = df.withColumn("is_boundary", F.when(F.col("value") != 1, 1).otherwise(0)) # 累加标记,每个边界点会开启一个新的子分组ID df = df.withColumn("sub_group_id", F.sum("is_boundary").over(main_window))
现在,name+sub_group_id就构成了子窗口的分区条件,你可以基于这个再定义子窗口:
sub_window = Window.partitionBy("name", "sub_group_id").orderBy("row_num")
这样就相当于在主窗口内部,又划分出了一个个独立的子窗口,用来处理更细粒度的计算。
问题2:替换连续1值为前一个非1值
针对你给出的DataFrame需求,咱们可以用上面的子窗口思路来实现,步骤很清晰:
第一步:先给数据加个行号,保证顺序
因为替换是按“当前非1值之后的连续1”来的,所以必须保证每个name分区内的行顺序和原始数据一致。如果你的数据没有自带顺序列,就加个行号:
# PySpark版本 df = df.withColumn("row_num", F.row_number().over(Window.partitionBy("name").orderBy(F.lit(1))))
第二步:标记边界,生成子分组ID
和问题1的思路一样,我们把每个非1值作为子分组的起始点,生成子分组ID:
df = df.withColumn("is_boundary", F.when(F.col("value") != 1, 1).otherwise(0)) main_window = Window.partitionBy("name").orderBy("row_num") df = df.withColumn("sub_group_id", F.sum("is_boundary").over(main_window))
第三步:在子分组内替换值
现在每个子分组里的第一个值就是非1值,我们只需要把整个子分组里的value都替换成这个第一个值就行:
sub_window = Window.partitionBy("name", "sub_group_id").orderBy("row_num") result_df = df.withColumn("value", F.first("value").over(sub_window)) # 最后去掉我们加的辅助列 result_df = result_df.drop("row_num", "is_boundary", "sub_group_id")
运行完之后,你得到的结果就和你期望的完全一致啦!
如果是用Pandas实现的话,逻辑更简洁:
import pandas as pd # 给每个name分区内的非1值生成分组ID df['sub_group_id'] = df.groupby('name')['value'].apply(lambda x: (x != 1).cumsum()) # 每个分组内的value都替换成分组的第一个值 result_df = df.assign(value=df.groupby(['name', 'sub_group_id'])['value'].transform('first')).drop('sub_group_id', axis=1)
内容的提问来源于stack exchange,提问作者meysam
相关产品推荐
相关产品推荐

