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

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 ID 13 会被误替换到 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

相关产品推荐
方舟 Agent Plan

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

最近更新时间:2026.06.23 11:16:06