PySpark需求:按行号延续父项值直至新父项出现
问题描述
现有如下结构的DataFrame:
| 订单编号 | 行号 | 物料编码 | 类型 |
|---|---|---|---|
| 12345 | 1 | 1001 | Parent |
| 12345 | 2 | 1002 | Child |
| 12345 | 3 | 1003 | Child |
| 12345 | 4 | 1004 | Child |
| 12345 | 5 | 1005 | Parent |
| 12345 | 6 | 1006 | Child |
需要新增一列Parent Item,规则为:
- 每个物料对应的父项是其之前最近的
Parent类型物料 - 父项自身的
Parent Item填自己的物料编码 - 父项编码持续填充直到出现新的父项
期望得到的结果:
| 行号 | 物料编码 | 类型 | Parent Item |
|---|---|---|---|
| 1 | 1001 | Parent | 1001 |
| 2 | 1002 | Child | 1001 |
| 3 | 1003 | Child | 1001 |
| 4 | 1004 | Child | 1001 |
| 5 | 1005 | Parent | 1005 |
| 6 | 1006 | Child | 1005 |
此前尝试过LAG函数和按订单+类型分区的窗口函数,均未实现需求。
解决方案
1. Pandas 实现
核心逻辑是标记父项位置后,用向前填充的方式继承最近的父项编码:
import pandas as pd # 构造示例数据 df = pd.DataFrame({ "订单编号": [12345]*6, "行号": [1,2,3,4,5,6], "物料编码": ["1001","1002","1003","1004","1005","1006"], "类型": ["Parent","Child","Child","Child","Parent","Child"] }) # 仅父项保留物料编码,子项设为缺失值 df["Parent Item"] = df.apply(lambda x: x["物料编码"] if x["类型"] == "Parent" else pd.NA, axis=1) # 向前填充缺失值,让子项继承最近的父项编码 df["Parent Item"] = df["Parent Item"].ffill() # 输出目标列 print(df[["行号","物料编码","类型","Parent Item"]])
2. PySpark 实现
针对大数据量场景,用窗口函数结合last()函数实现向前取值:
from pyspark.sql import SparkSession from pyspark.sql.window import Window from pyspark.sql.functions import when, last # 初始化Spark会话 spark = SparkSession.builder.appName("ParentItem").getOrCreate() # 构造示例数据 data = [ (12345,1,"1001","Parent"), (12345,2,"1002","Child"), (12345,3,"1003","Child"), (12345,4,"1004","Child"), (12345,5,"1005","Parent"), (12345,6,"1006","Child") ] df = spark.createDataFrame(data, ["订单编号","行号","物料编码","类型"]) # 定义窗口:按订单分区,按行号排序,范围从起始到当前行 window_spec = Window.partitionBy("订单编号").orderBy("行号").rowsBetween(Window.unboundedPreceding, Window.currentRow) # 提取父项编码并向前填充 df = df.withColumn( "Parent Item", last(when(df["类型"] == "Parent", df["物料编码"]), ignorenulls=True).over(window_spec) ) # 输出目标列 df.select("行号","物料编码","类型","Parent Item").show()
内容的提问来源于stack exchange,提问作者Chaddeus
相关产品推荐
相关产品推荐

