如何在PySpark中通过自定义函数逐行处理DataFrame并替换值?
在PySpark中实现指定的映射逻辑
PySpark不建议逐行遍历DataFrame(会严重影响性能),推荐用内置的矢量化操作来实现你的需求,以下是几种高效的实现方式:
方法一:使用when+otherwise(推荐)
利用PySpark SQL的条件判断函数,直接对列进行批量处理:
from pyspark.sql import functions as F df = df.withColumn( "new_col", F.when(F.col("temp") == "1 Person", "One") .when(F.col("temp") == "2 Persons", "Two") .when(F.col("temp") == "3 Persons", "Three") .when(F.col("temp").isin([ "4 Persons","5 Persons", "6 Persons", "7 Persons", "8 Persons","9 Persons","10 Persons","11 Persons" ]), "More") .otherwise(None) )
方法二:使用create_map构建映射表
先构建值到结果的映射,再通过coalesce匹配对应结果,未匹配到的返回None:
from pyspark.sql import functions as F from itertools import chain # 单个值对应结果的映射 single_mappings = [ ("1 Person", "One"), ("2 Persons", "Two"), ("3 Persons", "Three") ] # 多值统一映射到"More" more_values = [ "4 Persons","5 Persons", "6 Persons", "7 Persons", "8 Persons","9 Persons","10 Persons","11 Persons" ] more_mappings = [(val, "More") for val in more_values] # 合并所有映射并创建map对象 mapping = F.create_map(list(chain(*(single_mappings + more_mappings)))) # 应用映射,未匹配项返回None df = df.withColumn("new_col", F.coalesce(mapping[F.col("temp")], F.lit(None)))
方法三:使用UDF(不推荐)
如果一定要用类似Python自定义函数的方式,可以注册UDF,但UDF会脱离Spark的矢量化执行,性能远不如内置函数:
from pyspark.sql import functions as F from pyspark.sql.types import StringType def number(temp_val): if temp_val == '1 Person': return 'One' elif temp_val == '2 Persons': return 'Two' elif temp_val == '3 Persons': return 'Three' elif temp_val in ['4 Persons','5 Persons', '6 Persons', '7 Persons','8 Persons', '9 Persons','10 Persons','11 Persons']: return 'More' else: return None # 注册UDF number_udf = F.udf(number, StringType()) # 应用UDF到目标列 df = df.withColumn("new_col", number_udf(F.col("temp")))
注意事项
- 优先选择方法一或方法二,尤其是处理大数据量时,内置函数的执行效率远高于UDF。
- 如果
temp列有大量不同的值,方法二的映射表方式会更简洁。
内容的提问来源于stack exchange,提问作者ar_mm18
相关产品推荐
相关产品推荐

