如何在PySpark中基于字典为数据表新增映射列?
PySpark实现字典映射新增列
需求说明
基于给定的字典dictionary_Tag,为PySpark数据表新增description列,通过item_name字段匹配字典中的对应值,得到目标结果表。
原始字典
dictionary_Tag = {'A':'unitA&', 'B':'B&', 'C':'unitC', 'D':'D#' }
原始数据表
|item_name|item_value|timestamp |idx| +---------+----------+----------------------------+---+ |A |0.25 |2023-03-01T17:20:00.000+0000|0 | |B |0.34 |2023-03-01T17:20:00.000+0000|0 | |A |0.3 |2023-03-01T17:25:00.000+0000|1 | |B |0.54 |2023-03-01T17:25:00.000+0000|1 | |A |0.3 |2023-03-01T17:30:00.000+0000|2 | |B |0.54 |2023-03-01T17:30:00.000+0000|2 |
目标结果表
|item_name|item_value|timestamp |idx|description| +---------+----------+----------------------------+---+-----------+ |A |0.25 |2023-03-01T17:20:00.000+0000|0 |unitA& | |B |0.34 |2023-03-01T17:20:00.000+0000|0 |B& | |A |0.3 |2023-03-01T17:25:00.000+0000|1 |unitA& | |B |0.54 |2023-03-01T17:25:00.000+0000|1 |B& | |A |0.3 |2023-03-01T17:30:00.000+0000|2 |unitA& | |B |0.54 |2023-03-01T17:30:00.000+0000|2 |B& |
实现方法
方法一:使用create_map构建映射(推荐)
适合字典键值对较多的场景,代码简洁且可扩展性强。若字典数据量较大,可结合broadcast优化传输效率。
from pyspark.sql import SparkSession from pyspark.sql.functions import col, create_map, lit, broadcast # 初始化SparkSession(按需执行) spark = SparkSession.builder.appName("DictMapping").getOrCreate() # 将字典转为Spark映射表达式 mapping_expr = create_map(*[lit(x) for pair in dictionary_Tag.items() for x in pair]) # 新增description列 df_with_desc = df.withColumn("description", mapping_expr[col("item_name")]) # 大字典优化:广播映射表达式,减少节点间数据传输 # df_with_desc = df.withColumn("description", broadcast(mapping_expr)[col("item_name")]) # 查看结果 df_with_desc.show(truncate=False)
方法二:使用when+otherwise逐个匹配
适合字典键值对较少的场景,逻辑直观。
from pyspark.sql.functions import when, col # 逐个匹配字典键值,生成description列 df_with_desc = df.withColumn( "description", when(col("item_name") == "A", "unitA&") .when(col("item_name") == "B", "B&") .when(col("item_name") == "C", "unitC") .when(col("item_name") == "D", "D#") # 可选:为未匹配到的条目设置默认值 # .otherwise("Unknown") ) # 查看结果 df_with_desc.show(truncate=False)
内容的提问来源于stack exchange,提问作者MMV
相关产品推荐
相关产品推荐

