【问题标题】:How to Compare rows values in Pyspark using lead\lag?如何使用lead\\lag比较Pyspark中的行值?
【发布时间】:2022-11-12 12:19:14
【问题描述】:
我有一个列名称为“YEAR”的数据框,我想检查列的备用行是否匹配,如果备用值匹配,则更新另一个值为 100 的列“FLAG”。
df_prod
Year FLAG
2020 None
2020 None
2019 None
2021 None
2021 None
2022 None
预期产出
**
Year FLAG
2019 None
2020 None
2020 100
2021 None
2021 100
2022 None
**
【问题讨论】:
标签:
python
pyspark
databricks
【解决方案1】:
以下使用 Windowing 功能的 sn-p 应该为您执行此操作:
from pyspark.sql.window import Window
from pyspark.sql.functions import col, lag, when
df = spark.createDataFrame([(2020, None), (2020, None), (2019, None), (2021, None), (2021, None), (2022, None)], "Year: int, FLAG: int")
window = Window.partitionBy().orderBy("Year")
df.withColumn("FLAG", when(col("Year") == lag(col("Year")).over(window), 100)).show()
+----+----+
|Year|FLAG|
+----+----+
|2019|null|
|2020|null|
|2020| 100|
|2021|null|
|2021| 100|
|2022|null|
+----+----+