如何在PySpark中将字典值映射到DataFrame的新列中
错误原因
- 函数传参顺序颠倒:你定义的
my_mapp_fn第一个参数是待匹配的列名,第二个是映射字典,但你调用时传的第一个参数是字典、第二个是col('OrderId'),导致函数内部把Column对象当成字典处理,触发不可迭代报错。 - 变量名不匹配:函数内部循环用了
d.items(),但你的参数名是dict1,运行时会报变量未定义错误。 - 列名大小写不匹配:你DataFrame里的列名是
OrderID,调用时写的是OrderId,后续可能触发列不存在报错。 - 变量名不规范:用Python内置关键字
dict作为自定义变量名,容易引发命名冲突。 - 映射逻辑与需求不匹配:原有函数返回的是字典对应的数字值,你需要的是
Y/N标识,需要额外补充判断逻辑。
修复后的代码
首先导入依赖:
from pyspark.sql.functions import col, when, lit, coalesce
方案1:直接匹配需求(适合固定规则场景)
从你给出的示例可以看出,仅OrderID=667593514时返回N,其余返回Y,可以直接简化逻辑:
# 读取数据 df = spark.read.option("header", True).csv("sample.csv") # 新增SomeFlag列 new_df = df.withColumn("SomeFlag", when(col("OrderID") == "667593514", lit("N")).otherwise(lit("Y"))) # 验证结果 new_df.select("SomeFlag").show()
方案2:基于字典映射(适合后续字典会扩展的场景)
如果后续映射规则会变化,需要保留字典映射逻辑,可以使用修复后的代码:
# 避免用dict作为变量名 order_id_map = {'443368995': 0, '667593514': 1, '940995585': 2, '880811536': 3, '174590194': 4} def my_map_fn(check_col_name, map_dict): # 匹配字典键返回对应值,匹配不到返回-1 return coalesce(*[when(col(check_col_name) == key, lit(value)) for key, value in map_dict.items()], lit(-1)) df = spark.read.option("header", True).csv("sample.csv") # 先映射得到字典对应值 df_with_map = df.withColumn("map_value", my_map_fn("OrderID", order_id_map)) # 按规则生成Y/N标识 new_df = df_with_map.withColumn("SomeFlag", when(col("map_value") == 1, lit("N")).otherwise(lit("Y"))) # 验证结果 new_df.select("SomeFlag").show()
运行后即可得到你期望的SomeFlag列结果。
内容的提问来源于stack exchange,提问作者Varun Singh
相关产品推荐
相关产品推荐

