如何使用元组作为字典键替换PySpark DataFrame中的值
PySpark 统一银行名称映射实现方案
问题场景
现有包含银行名称列的PySpark DataFrame,同一银行存在多种名称表述,示例数据如下:
+---+--------------------+ | id| name| +---+--------------------+ | 1| BANCO SANTANDER| | 2| SANTANDER| | 3|BANCO SANTANDER S.A.| | 4| JP MORGAN| | 5| JP MORGAN CHASE| | 6| CITIBANK| | 7| CITI| | 8| CITIGROUP| | 9| HSBC HOLDINGS| | 10| HBSC| +---+--------------------+
为避免编写大量CASE WHEN语句,预先定义了以元组为键的映射字典,多个别名对应同一标准名称:
bank_dict = { ('JP MORGAN CHASE',):'JP MORGAN', ('CITI', 'CITIGROUP'):'CITIBANK', ('BANCO SANTANDER', 'BANCO SANTANDER S.A.', 'SANTANDER CREDIT CARDS'):'SANTANDER', ('HSBC HOLDINGS',):'HSBC' }
需求:将DataFrame中name列匹配字典键中任意元素的值,替换为对应的标准名称,未匹配的保留原名称,预期结果如下:
+---+--------------------+---------+ | id| name| new_name| +---+--------------------+---------+ | 1| BANCO SANTANDER|SANTANDER| | 2| SANTANDER|SANTANDER| | 3|BANCO SANTANDER S.A.|SANTANDER| | 4| JP MORGAN|JP MORGAN| | 5| JP MORGAN CHASE|JP MORGAN| | 6| CITIBANK| CITIBANK| | 7| CITI| CITIBANK| | 8| CITIGROUP| CITIBANK| | 9| HSBC HOLDINGS| HSBC| | 10| HBSC| HBSC| +---+--------------------+---------+
注:原预期结果中第9行的HBSC应为笔误,按照字典映射逻辑,HSBC HOLDINGS应转换为HSBC
实现步骤
1. 转换映射字典结构
原字典是「多别名→单标准名」的结构,我们需要将其转换为「单别名→单标准名」的扁平结构,同时补充标准名到自身的映射,确保标准名本身不会被当成未匹配项:
# 构建扁平映射字典 name_mapping = {} for aliases, standard_name in bank_dict.items(): # 遍历每个别名,添加到映射 for alias in aliases: name_mapping[alias] = standard_name # 补充标准名自身的映射(避免标准名因不在别名列表中被保留原值,逻辑上更严谨) standard_names = set(bank_dict.values()) for std_name in standard_names: if std_name not in name_mapping: name_mapping[std_name] = std_name
2. PySpark中实现映射转换
利用PySpark的create_map构建映射表达式,结合coalesce实现「匹配则替换,否则保留原值」的逻辑:
from pyspark.sql import SparkSession from pyspark.sql.functions import col, create_map, lit, coalesce # 初始化SparkSession spark = SparkSession.builder.appName("BankNameStandardization").getOrCreate() # 加载示例数据(实际场景替换为你的数据源) data = [ (1, "BANCO SANTANDER"), (2, "SANTANDER"), (3, "BANCO SANTANDER S.A."), (4, "JP MORGAN"), (5, "JP MORGAN CHASE"), (6, "CITIBANK"), (7, "CITI"), (8, "CITIGROUP"), (9, "HSBC HOLDINGS"), (10, "HBSC") ] df = spark.createDataFrame(data, ["id", "name"]) # 构建Spark映射表达式 map_expr = create_map( *[lit(item) for pair in name_mapping.items() for item in pair] ) # 生成new_name列 df_result = df.withColumn( "new_name", coalesce(map_expr[col("name")], col("name")) ) # 查看结果 df_result.show(truncate=False)
方案优势
- 避免了冗长的
CASE WHEN链式判断,代码更简洁易维护 - 基于Spark的
map操作实现O(1)查找,处理大规模数据时性能更优 - 映射字典的结构易于扩展,新增银行别名只需修改原字典即可
内容的提问来源于stack exchange,提问作者Marlon Iwanaga
相关产品推荐
相关产品推荐

