PySpark中基于字典替换DataFrame FORMULA列值生成新列的问题
问题描述
现有一个PySpark DataFrame df_1,结构如下:
+------+-------------+----------+-----+ |KPI_ID| FORMULA|PARAMETERS|VALUE| +------+-------------+----------+-----+ | 132|131 * (1-133)| 131,133| null| | 135| 134 * 1000| 134| null| | 138| 137 - 139| 137,139| null| | 143| 142 * 137| 142,137| null| +------+-------------+----------+-----+
另有存储KPI对应值的字典 kpi_values:
{'123': 'None', '124': '763.675293', '125': '1249.6796875', '126': '9313.3896484', '127': '14742.0', '128': 'None', '129': '2.8590682', '130': '6010.1899414', '131': '2.96', '133': '83.3741302', '134': '-1.0', '136': '955.4359741', '137': '27989.8808594', '139': '1674.25', '140': '97000.0', '141': '2388.3144531', '142': '0.0', '144': '7421.25', '145': '0.0', '146': '1.3425742'}
需要生成包含 FORMULAREPLACED 列的输出DataFrame,该列将 FORMULA 中的KPI ID替换为字典中的对应值,预期输出如下:
+------+-------------+----------+-----+-------------------------+ |KPI_ID| FORMULA|PARAMETERS|VALUE|FORMULAREPLACED | +------+-------------+----------+-----+-------------------------+ | 132|131 * (1-133)| 131,133| null| 2.96 * (1-83.3741302) | | 135| 134 * 1000| 134| null| -1.0 * 1000 | | 138| 137 - 139| 137,139| null| 27989.8808594 - 1674.25| | 143| 142 * 137| 142,137| null| 0.0 * 27989.8808594 | +------+-------------+----------+-----+-------------------------+
原尝试代码未得到预期结果:
for kpi_id, value in kpi_values.items(): placeholder = f"{kpi_id}" if value is None: value = '0' else: value # Handle NULL values df_1 = df_1.withColumn("FORMULAREPLACED",F.expr(f"regexp_replace(FORMULA,'{placeholder }','{value}')")) df_1 = df_1.withColumn("FORMULA", F.expr(f"replace(FORMULA, '{placeholder }', '{value}')")) df_1=df_1.replace(list(kpi_values.keys()), list(kpi_values.values()), 'FORMULA')
限制条件:Unity Catalog不支持UDF,需用其他方法实现。
解决方案
原代码问题分析
- 循环中每次调用
withColumn都会覆盖之前的FORMULAREPLACED列,最终仅保留最后一次替换结果 replace方法会替换所有匹配子串,易出现误替换(比如KPI ID13会被误替换到131中)- 字典中值为
'None'的处理逻辑错误,原代码判断value is None,但实际是字符串'None'
正确实现方法
通过构建链式 regexp_replace 表达式,结合单词边界确保只替换完整KPI ID,避免误匹配:
from pyspark.sql import functions as F # 预处理字典:将字符串'None'替换为'0' processed_kpi = {k: '0' if v == 'None' else v for k, v in kpi_values.items()} # 构建链式替换表达式,初始为FORMULA列,依次替换每个KPI ID replace_expr = "FORMULA" for kpi_id, val in processed_kpi.items(): # 使用\b匹配单词边界,保证只替换独立的KPI ID replace_expr = f"regexp_replace({replace_expr}, '\\b{kpi_id}\\b', '{val}')" # 添加FORMULAREPLACED列 df_result = df_1.withColumn("FORMULAREPLACED", F.expr(replace_expr)) # 查看结果 df_result.show(truncate=False)
说明
- 单词边界
\b避免了短KPI ID误替换长ID的问题 - 先预处理字典,统一处理
'None'值为'0' - 链式构建替换表达式,一次性完成所有替换,避免循环覆盖列的问题
- 完全使用Spark内置函数实现,兼容Unity Catalog
内容的提问来源于stack exchange,提问作者sakura angadi
相关产品推荐
相关产品推荐

