PySpark需求:为同ID的DataFrame添加prev_value列以计算值差异
嘿,作为PySpark新手碰到这种需求完全没问题,我来一步步教你实现!
首先得明确:PySpark的DataFrame本身是无序的,所以我们得先给数据加上一个能确定顺序的标识,不然lag(取上一行值)的结果会乱掉。
实现你给出的期望结果(全局上一行值)
按照你提供的输出示例,prev_value取的是整个数据集的上一行值(哪怕id不同),具体步骤如下:
- 导入需要的函数和窗口类:
from pyspark.sql import Window from pyspark.sql.functions import lag, monotonically_increasing_id
- 给原DataFrame添加一个自增行号,用来固定数据的顺序:
# 添加行号列,确保数据和你给出的输入顺序一致 df_with_row = df.withColumn("row_id", monotonically_increasing_id())
- 定义全局排序的窗口规则:
# 按行号排序,保证取到的是全局的上一行 global_window = Window.orderBy("row_id")
- 使用
lag函数生成prev_value列,最后删掉临时的行号列:
result_df = df_with_row.withColumn("prev_value", lag("value").over(global_window)) \ .drop("row_id") # 查看结果 result_df.show()
运行后就能得到你想要的输出:
+---+-----+----------+ | id|value|prev_value| +---+-----+----------+ | 1| 65| null| | 1| 66| 65| | 1| 65| 66| | 2| 68| 65| | 2| 71| 68| +---+-----+----------+
补充:如果是同一id内取上一行值(更符合"同一id下差值计算"的需求)
如果你其实是想让每个id组内的第一行prev_value为null,组内其他行取组内上一行的值(这样计算同一id的value差值更合理),只需要修改窗口规则,加上partitionBy("id")即可:
# 按id分组,组内按行号排序 group_window = Window.partitionBy("id").orderBy("row_id") result_df_group = df_with_row.withColumn("prev_value", lag("value").over(group_window)) \ .drop("row_id") result_df_group.show()
这时候输出会是:
+---+-----+----------+ | id|value|prev_value| +---+-----+----------+ | 1| 65| null| | 1| 66| 65| | 1| 65| 66| | 2| 68| null| | 2| 71| 68| +---+-----+----------+
这样后续计算同一id下的value差值,直接用value - prev_value就可以啦!
内容的提问来源于stack exchange,提问作者Myat Noe
相关产品推荐
相关产品推荐

