如何用PySpark窗口函数实现指定DataFrame的win_name批量更新?
能否用PySpark窗口函数实现指定DataFrame更新需求?
输入DataFrame
| name | type | description | win_name | value |
|---|---|---|---|---|
| sonu | doctor | dentist | monu | false |
| monu | doctor | ortho | monu | true |
| chintu | doctor | dentist | monu | false |
目标DataFrame(修改sonu的value为true后,所有行win_name设为sonu)
| name | type | description | win_name | value |
|---|---|---|---|---|
| sonu | doctor | dentist | sonu | true |
| monu | doctor | ortho | sonu | false |
| chintu | doctor | dentist | sonu | false |
实现方案:可以通过PySpark窗口函数完成
核心思路是:先更新sonu的value字段为true,再利用全表窗口提取符合条件(name='sonu'且value=true)的name值,将其统一赋值给所有行的win_name字段。
完整代码示例
from pyspark.sql import SparkSession from pyspark.sql.window import Window import pyspark.sql.functions as F # 初始化SparkSession spark = SparkSession.builder.appName("win_name_update").getOrCreate() # 创建输入DataFrame data = [ ("sonu", "doctor", "dentist", "monu", False), ("monu", "doctor", "ortho", "monu", True), ("chintu", "doctor", "dentist", "monu", False) ] input_df = spark.createDataFrame(data, ["name", "type", "description", "win_name", "value"]) # 1. 更新sonu的value为true updated_df = input_df.withColumn( "value", F.when(F.col("name") == "sonu", True).otherwise(F.col("value")) ) # 2. 定义全表窗口(覆盖所有行,无分区、无排序) full_window = Window.rowsBetween(Window.unboundedPreceding, Window.unboundedFollowing) # 3. 提取sonu的name作为所有行的win_name result_df = updated_df.withColumn( "win_name", # 取符合条件的第一个name值,ignorenulls确保不会因为空值报错 F.first( F.when(F.col("name") == "sonu" & F.col("value") == True, F.col("name")), ignorenulls=True ).over(full_window) ) # 查看结果 result_df.show()
代码解释
- 更新value字段:使用
when-otherwise条件表达式,将name='sonu'的行的value设为true,其他行保持原value。 - 全表窗口定义:
Window.rowsBetween(Window.unboundedPreceding, Window.unboundedFollowing)表示窗口包含DataFrame的所有行,实现跨行共享数据。 - 统一赋值win_name:通过
first函数结合when条件,筛选出name='sonu'且value=true的行的name值,然后将该值应用到窗口内的所有行,实现win_name的统一更新。
替代实现方式
也可以使用collect_set结合element_at函数实现,效果一致:
result_df = updated_df.withColumn( "win_name", F.element_at( F.collect_set(F.when(F.col("name") == "sonu" & F.col("value") == True, F.col("name"))).over(full_window), 1 ) )
内容的提问来源于stack exchange,提问作者Rishabh Soni
相关产品推荐
相关产品推荐

