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

如何在PySpark中计算DataFrame内相互依赖的amt_c、lkp_rev列

PySpark实现递归计算字段方案

你的需求属于分组内的递归关联计算,amt_c和lkp_rev存在逐行依赖关系,普通窗口函数无法直接实现,我们可以用applyInPandas按分组处理,逻辑清晰易维护。


完整实现代码

import pyspark.sql.functions as f
from pyspark.sql import Window
from pyspark.sql.types import StructType, StructField, IntegerType

# 原始DataFrame
df = spark.createDataFrame([
(124,10,8),
(124,20,7),
(125,30,6),
(125,40,5),
(126,50,4),
(126,60,3),
(126,70,2),
(127,80,1)],("ACC_KEY", "AMT", "value"))

# 1. 生成分组行号row_now
window_spec = Window.partitionBy("ACC_KEY").orderBy(f.col('AMT').desc())
df_with_row = df.withColumn("row_now", f.row_number().over(window_spec))

# 2. 定义输出结果的schema
output_schema = StructType([
    StructField("ACC_KEY", IntegerType(), nullable=False),
    StructField("AMT", IntegerType(), nullable=False),
    StructField("value", IntegerType(), nullable=False),
    StructField("row_now", IntegerType(), nullable=False),
    StructField("amt_c", IntegerType(), nullable=False),
    StructField("lkp_rev", IntegerType(), nullable=False)
])

# 3. 定义单分组计算逻辑
def calc_single_group(pdf):
    # 入参pdf是单个ACC_KEY对应的pandas DataFrame,已按AMT降序排序
    amt_c_list = []
    lkp_rev_list = []
    pre_lkp_rev = None
    for _, row in pdf.iterrows():
        # 计算amt_c
        if row["row_now"] == 1:
            current_amt_c = row["AMT"]
        else:
            current_amt_c = pre_lkp_rev
        # 计算lkp_rev
        current_lkp_rev = current_amt_c + row["value"]
        # 存入结果列表
        amt_c_list.append(current_amt_c)
        lkp_rev_list.append(current_lkp_rev)
        pre_lkp_rev = current_lkp_rev
    pdf["amt_c"] = amt_c_list
    pdf["lkp_rev"] = lkp_rev_list
    return pdf

# 4. 分组应用计算逻辑
result_df = df_with_row.groupBy("ACC_KEY").applyInPandas(calc_single_group, schema=output_schema)

# 按期望的顺序输出验证
result_df.orderBy(f.col("ACC_KEY").desc(), f.col("row_now")).show()

输出结果

和你要求的预期结果完全一致:

+-------+---+-----+-------+-----+-------+
|ACC_KEY|AMT|value|row_now|amt_c|lkp_rev|
+-------+---+-----+-------+-----+-------+
|    126| 70|    2|      1|   70|     72|
|    126| 60|    3|      2|   72|     75|
|    126| 50|    4|      3|   75|     79|
|    124| 20|    7|      1|   20|     27|
|    124| 10|    8|      2|   27|     35|
|    127| 80|    1|      1|   80|     81|
|    125| 40|    5|      1|   40|     45|
|    125| 30|    6|      2|   45|     51|
+-------+---+-----+-------+-----+-------+

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

相关产品推荐
方舟 Agent Plan

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

最近更新时间:2026.09.29 11:27:03