You need to enable JavaScript to run this app.
优惠活动
大模型
产品
解决方案
定价
更多

PySpark:基于条件匹配同SUBS_NO行更新DataFrame列值

PySpark实现指定条件下的列值替换需求

需求说明

当DataFrame中某行的RELATION字段不等于'1'时,将该行的LAST_NAME、FIRST_NAME、BILL_NO、ADDRESS字段,替换为同SUBS_NO且RELATION为'1'的行的对应值。

解决方案代码

from pyspark.sql import SparkSession
from pyspark.sql.functions import col, when

# 初始化SparkSession
spark = SparkSession.builder.appName("ReplaceColumns").getOrCreate()

# 创建原始测试数据
data = [
    (1, "SMITH", "JIM", "123", "123", "123 MAIN ST", "SOMEWHERE", "WI", "12345", "1"),
    (2, "DOE", "SALLY", "456", "456", "456 ELM AVE", "ANYWHERE", "AZ", "54321", "1"),
    (3, "SMITH", "JANE", "123", "789", "888 3RD ST", "SOMEWHERE", "WI", "12345", "2")
]

columns = ["ID", "LAST_NAME", "FIRST_NAME", "SUBS_NO", "BILL_NO", "ADDRESS", "CITY", "STATE", "ZIP", "RELATION"]
df = spark.createDataFrame(data, columns)

# 步骤1:提取每个SUBS_NO对应的RELATION='1'的基准数据
base_df = df.filter(col("RELATION") == "1") \
            .select("SUBS_NO", "LAST_NAME", "FIRST_NAME", "BILL_NO", "ADDRESS") \
            .withColumnRenamed("LAST_NAME", "BASE_LAST_NAME") \
            .withColumnRenamed("FIRST_NAME", "BASE_FIRST_NAME") \
            .withColumnRenamed("BILL_NO", "BASE_BILL_NO") \
            .withColumnRenamed("ADDRESS", "BASE_ADDRESS")

# 步骤2:关联原表和基准数据,保留所有原始行
joined_df = df.join(base_df, on="SUBS_NO", how="left")

# 步骤3:根据条件替换指定列值
result_df = joined_df.withColumn("LAST_NAME", 
                                 when(col("RELATION") != "1", col("BASE_LAST_NAME")).otherwise(col("LAST_NAME"))) \
                     .withColumn("FIRST_NAME", 
                                 when(col("RELATION") != "1", col("BASE_FIRST_NAME")).otherwise(col("FIRST_NAME"))) \
                     .withColumn("BILL_NO", 
                                 when(col("RELATION") != "1", col("BASE_BILL_NO")).otherwise(col("BILL_NO"))) \
                     .withColumn("ADDRESS", 
                                 when(col("RELATION") != "1", col("BASE_ADDRESS")).otherwise(col("ADDRESS"))) \
                     .drop("BASE_LAST_NAME", "BASE_FIRST_NAME", "BASE_BILL_NO", "BASE_ADDRESS")

# 查看结果
result_df.show()

代码解释

  • 提取基准数据:先过滤出RELATION='1'的行,只保留需要用来替换的列,并给这些列加上前缀BASE_,避免关联后列名冲突。
  • 关联数据:用left join将原表和基准数据按SUBS_NO关联,确保所有原始行都被保留,即使某个SUBS_NO没有RELATION='1'的行(这种情况下替换列会保持原值)。
  • 条件替换:使用when函数判断RELATION是否不等于'1',如果是则用基准列的值替换原列,否则保留原列值,最后删除临时的基准列。

输出结果

执行代码后会得到符合预期的DataFrame:

+---+---------+----------+-------+-------+-------------+---------+-----+-----+--------+
| ID|LAST_NAME|FIRST_NAME|SUBS_NO|BILL_NO|       ADDRESS|     CITY|STATE|  ZIP|RELATION|
+---+---------+----------+-------+-------+-------------+---------+-----+-----+--------+
|  1|    SMITH|        JIM|    123|    123| 123 MAIN ST |SOMEWHERE|   WI|12345|       1|
|  2|      DOE|      SALLY|    456|    456|456 ELM AVE  | ANYWHERE|   AZ|54321|       1|
|  3|    SMITH|        JIM|    123|    123| 123 MAIN ST |SOMEWHERE|   WI|12345|       2|
+---+---------+----------+-------+-------+-------------+---------+-----+-----+--------+

内容的提问来源于stack exchange,提问作者Craig Larson

相关产品推荐
方舟 Agent Plan

超全模态模型 × Harness 升级,最新支持 Deepseek-V4.1-Flash、GLM-5.3 系列、Doubao-Seedream-5.0-pro、Kimi-K3 (部分), 限时 9.9 元起

最近更新时间:2026.07.04 17:28:19