如何基于字典值为Spark DataFrame创建新列?解决TypeError报错
解决Spark DataFrame基于汇率字典添加转换列的问题
错误原因
你触发的TypeError: unhashable type: 'Column',是因为raw_df.currency是Spark的Column对象,而Python字典要求key必须是可哈希的(比如字符串、数字等),Column对象不满足这个条件,所以不能直接用它作为字典的key来取值。
两种可行解决方案
方案1:使用when-otherwise逐个匹配货币(适合货币种类少的场景)
通过Spark的条件函数when,针对每种货币分别设置汇率计算逻辑:
from pyspark.sql import functions as F # 原DataFrame raw_df = spark.createDataFrame([("USD", 1.00), ("EUR", 2.00)], ["currency", "value"]) # 添加转换列 raw_df = raw_df.withColumn( "value_eur", F.when(F.col("currency") == "USD", F.col("value") * 0.90) .when(F.col("currency") == "EUR", F.col("value") * 1.00) .otherwise(F.col("value")) # 处理未定义的货币,可根据需求修改默认值 ) raw_df.show()
输出结果:
+---------+-----+---------+ |currency|value|value_eur| +---------+-----+---------+ | USD| 1.0| 0.9| | EUR| 2.0| 2.0| +---------+-----+---------+
方案2:将字典转为Spark MapType列(适合货币种类多的场景)
把Python汇率字典转换成Spark支持的MapType常量列,再通过列的getItem方法动态获取对应汇率,这种方式不需要修改代码,只需更新字典即可适配新的货币:
from pyspark.sql import functions as F # 原汇率字典 currency_exchanges = {"EUR": 1.00, "USD": 0.90} # 将字典转为Spark MapType列 exchange_map = F.create_map(*[F.lit(item) for pair in currency_exchanges.items() for item in pair]) # 原DataFrame raw_df = spark.createDataFrame([("USD", 1.00), ("EUR", 2.00)], ["currency", "value"]) # 添加转换列,若货币不在字典中则返回null,可通过coalesce设置默认值 raw_df = raw_df.withColumn( "value_eur", F.col("value") * F.coalesce(exchange_map.getItem(F.col("currency")), F.lit(1.0)) ) raw_df.show()
输出结果和方案1一致,若有未定义的货币会使用默认汇率1.0,你可以根据需求调整coalesce里的默认值。
内容的提问来源于stack exchange,提问作者Andrea Nerla
相关产品推荐
相关产品推荐

