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
相关产品推荐
相关产品推荐

