【问题标题】:PySpark window function improvementPySpark 窗口功能改进
【发布时间】:2021-05-06 03:35:55
【问题描述】:

我需要用以前的记录值替换,所以我使用窗口函数实现了这个,但我想提高性能。您能否告知是否有其他替代方法。

from pyspark.sql import SparkSession, Window, DataFrame
from pyspark.sql.types import *
from pyspark.sql import functions as F

source = [(1,2,3),(2,3,4),(1,3,4)]
target = [(1,3,1),(3,4,1)]
schema = ['key','col1','col2']
source_df = spark.createDataFrame(source, schema=schema)
target_df = spark.createDataFrame(source, schema=schema)

df = source_df.unionAll(target_df)

window = Window.partitionBy(F.col('key')).orderBy(F.col('col2').asc())


df = df.withColumn('col1_prev', F.lag(F.col('col1_start')).over(window)\
       .withColumn('col1', F.lit('col1_next'))

df.show()

1,3,1
1,2,1
1,3,3
2,3,4
3,4,1

【问题讨论】:

    标签: python pyspark hive window-functions


    【解决方案1】:

    您可以在指定的时间间隔内使用last 函数,例如窗口中的最后两行。我这里以maxsize为例:

    import sys
    window = Window.partitionBy('key')\
                   .orderBy('col2')\
                   .rowsBetween(-sys.maxsize, -1)
    
    df = F.last(df['col1_prev'], ignorenulls=True).over(window)
    

    希望它能解决你的问题。

    【讨论】:

    • 谢谢 Amir,这是为了提高性能还是避免空值?我认为 last 将根据排序顺序返回最后一个值,但我需要以前的记录非空值
    • last 函数根据您想要的顺序返回最后一条记录(上一条记录)。您可以忽略或不忽略空值,这是您手中的一个很好的选择。我认为与lag 函数相比,last 应该更快。
    猜你喜欢
    • 2019-02-15
    • 1970-01-01
    • 1970-01-01
    • 1970-01-01
    • 1970-01-01
    • 1970-01-01
    • 1970-01-01
    • 2019-02-16
    • 1970-01-01
    相关资源
    最近更新 更多