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

Apache Spark(PySpark):如何用同列其他行值替换指定行列值

问题:替换指定行的列值失败

需求说明

我想要把LAST_NAME为'Maltster'的行中,ANNUAL_HOUSEHOLD_INCOME列的值替换为LAST_NAME为'Attiwill'的行中同一列的值。

数据示例

执行代码前的数据表:

+---------+-----------------------+
|LAST_NAME|ANNUAL_HOUSEHOLD_INCOME|                               
+---------+-----------------------+
|Maltster |20000                  |
|Attiwill |100000                 |
+---------+-----------------------+

期望执行后的数据表:

+---------+-----------------------+
|LAST_NAME|ANNUAL_HOUSEHOLD_INCOME|                               
+---------+-----------------------+
|Maltster |100000                 |
|Attiwill |100000                 |
+---------+-----------------------+

问题代码

运行以下代码后,目标值并未被覆盖:

df1.filter(F.col('LAST_NAME') == 'Maltster').withColumn("ANNUAL_HOUSEHOLD_INCOME", df1.filter(F.col('LAST_NAME') == 'Attiwill').select(F.col('ANNUAL_HOUSEHOLD_INCOME'))[0]).show()

问题原因与解决方案

原因分析

这段代码只对过滤出的Maltster子DataFrame做了修改,但原DataFramedf1没有被更新,也没有保存修改后的结果。另外,直接取select结果的[0]不是正确获取具体数值的方式,无法正确赋值。

正确实现方式

方法1:条件赋值(简单场景首选)

直接在原DataFrame上通过when/otherwise判断替换:

from pyspark.sql import functions as F

# 提取Attiwill的收入数值
attiwill_income = df1.filter(F.col('LAST_NAME') == 'Attiwill').select('ANNUAL_HOUSEHOLD_INCOME').first()[0]

# 对指定行替换列值
updated_df = df1.withColumn(
    "ANNUAL_HOUSEHOLD_INCOME",
    F.when(F.col('LAST_NAME') == 'Maltster', attiwill_income).otherwise(F.col('ANNUAL_HOUSEHOLD_INCOME'))
)

updated_df.show()

方法2:DataFrame合并(适合批量替换场景)

如果需要替换的行较多,可以通过关联方式批量更新:

from pyspark.sql import functions as F

# 提取目标替换值并构建匹配条件
attiwill_df = df1.filter(F.col('LAST_NAME') == 'Attiwill').select(
    F.lit('Maltster').alias('TARGET_NAME'),
    F.col('ANNUAL_HOUSEHOLD_INCOME').alias('NEW_INCOME')
)

# 关联后更新列值
updated_df = df1.join(attiwill_df, df1.LAST_NAME == attiwill_df.TARGET_NAME, 'left') \
    .withColumn(
        "ANNUAL_HOUSEHOLD_INCOME",
        F.when(F.col('TARGET_NAME').isNotNull(), F.col('NEW_INCOME')).otherwise(F.col('ANNUAL_HOUSEHOLD_INCOME'))
    ) \
    .drop('TARGET_NAME', 'NEW_INCOME')

updated_df.show()

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

相关产品推荐
方舟 Agent Plan

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

最近更新时间:2026.08.22 15:03:53