You need to enable JavaScript to run this app.
优惠活动
大模型
产品
解决方案
定价
更多

如何使用元组作为字典键替换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

相关产品推荐
方舟 Agent Plan

超全模态模型 × Harness 升级,最新支持 Deepseek-V4.1-Flash、GLM-5.3 系列、Doubao-Seedream-5.0-pro、Kimi-K3 (部分), 限时 9.9 元起

最近更新时间:2026.08.16 08:00:57