在PySpark中不使用withColumn如何基于另一列值更新DataFrame列
Spark DataFrame 直接更新已有列的实现方案
首先需要明确一个Spark API的常用特性:withColumn本身就支持直接覆盖已存在的列,不需要新建列再替换的迂回操作。当你传入的列名参数是当前DataFrame中已有的列名时,新生成的列会直接替换原有同名列的内容。
针对你的需求(仅满足status条件时更新req_id,其余场景保留req_id原有值),可以直接用如下原生DataFrame API实现:
val updatedDf = df.withColumn("req_id", when(col("status") === "9088", "Generated") .when(col("status") === "9089", "Deactive") .otherwise(col("req_id")) // 不匹配条件时直接返回原有req_id的取值 )
补充说明
- 该写法的执行逻辑和你用
spark.sql执行更新查询的逻辑完全等价,执行效率没有差异,全程都是Spark原生API操作。 - 如果你更偏好
select写法,也可以通过通配符+同名列覆盖实现相同效果,逻辑和上述写法一致:
val updatedDf = df.select( col("*"), when(col("status") === "9088", "Generated") .when(col("status") === "9089", "Deactive") .otherwise(col("req_id")) .alias("req_id") )
内容的提问来源于stack exchange,提问作者Metadata
相关产品推荐
相关产品推荐

