如何在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
相关产品推荐
相关产品推荐

