如何在PySpark中实现类似SparkR示例的窗口分区lag操作
SparkR 窗口 lag 逻辑的 PySpark 实现
以下是和你提供的SparkR逻辑完全对齐的PySpark实现代码:
# 导入所需依赖 from pyspark.sql import SparkSession from pyspark.sql.functions import lag, col from pyspark.sql.window import Window # 初始化SparkSession spark = SparkSession.builder.appName("lag_demo").getOrCreate() # 1. 对应SparkR中createDataFrame(mtcars)的步骤 # 你可以直接读取自己的mtcars数据源,这里给出模拟数据示例可直接运行 mtcars_sample = [ (21.0, 110, 1), (21.0, 110, 1), (22.8, 93, 1), (21.4, 110, 0), (18.7, 175, 0), (18.1, 105, 0), (14.3, 245, 0), (24.4, 62, 1) ] df = spark.createDataFrame(mtcars_sample, schema=["mpg", "hp", "am"]) # 2. 对应SparkR的窗口定义:按am字段分区,按hp字段排序 ws = Window.partitionBy("am").orderBy("hp") # 3. 对应SparkR的lag计算+列选择逻辑 out = df.select( lag(col("mpg"), 1).over(ws).alias("lag_mpg"), # 计算滞后1行的mpg值,指定别名方便后续使用 col("mpg"), col("hp"), col("am") ) # 验证结果 out.show()
逻辑对应说明
- PySpark的窗口规则通过
Window类定义,partitionBy对应SparkR的windowPartitionBy方法,链式调用orderBy对应SparkR的orderBy参数配置 lag函数第二个参数为偏移行数,传1和原代码逻辑一致,你也可以根据需求调整偏移量,添加第三个参数可以自定义分区内首行的默认lag值
内容的提问来源于stack exchange,提问作者Innovator-programmer
相关产品推荐
相关产品推荐

