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

Spark DataFrame中利用字典推导式创建新旧值对比列的方法

解决方案:Spark DataFrame 生成新旧值差异列

方法1:使用Python UDF实现

如果你的sot和metrics列是Map类型或者JSON字符串,可以先将其转为字典,通过字典推导式对比差异后封装成UDF使用。

步骤示例:

  1. 定义差异计算逻辑:对比两个字典的键值对,仅保留值不同的键,生成{key: {"old": 旧值, "new": 新值}}格式的结果。
  2. 注册UDF并应用到DataFrame。
from pyspark.sql import SparkSession
from pyspark.sql.functions import udf
from pyspark.sql.types import MapType, StringType, StructType, StructField

# 初始化SparkSession
spark = SparkSession.builder.appName("diff_column").getOrCreate()

# 示例数据:sot和metrics为Map[String, String]类型
data = [
    ({"a": "100", "b": "200", "c": "300"}, {"a": "100", "b": "250", "c": "300"}),
    ({"x": "50", "y": "60"}, {"x": "55", "y": "60"})
]
df = spark.createDataFrame(data, ["sot", "metrics"])

# 定义差异计算函数
def calculate_diff(sot_dict, metrics_dict):
    diff = {}
    # 仅保留双方都存在且值不同的键
    for key in sot_dict:
        if key in metrics_dict and sot_dict[key] != metrics_dict[key]:
            diff[key] = {"old": sot_dict[key], "new": metrics_dict[key]}
    return diff

# 定义UDF返回类型:Map[String, Struct(old: String, new: String)]
diff_schema = MapType(
    StringType(),
    StructType([
        StructField("old", StringType()),
        StructField("new", StringType())
    ])
)
diff_udf = udf(calculate_diff, diff_schema)

# 生成差异列
result_df = df.withColumn("value_diff", diff_udf("sot", "metrics"))
result_df.show(truncate=False)

输出结果:

+-----------------------------+-----------------------------+-----------------------+
|sot                          |metrics                      |value_diff             |
+-----------------------------+-----------------------------+-----------------------+
|{a -> 100, b -> 200, c -> 300}|{a -> 100, b -> 250, c -> 300}|{b -> {old:200, new:250}}|
|{x -> 50, y -> 60}           |{x -> 55, y -> 60}           |{x -> {old:50, new:55}}  |
+-----------------------------+-----------------------------+-----------------------+

方法2:使用Spark内置函数(无UDF,性能更优)

若你的Spark版本在3.0及以上,可借助内置高阶函数(map_entries、filter、transform、map_from_entries)实现,避免UDF带来的Python-JVM序列化开销。

from pyspark.sql.functions import map_entries, filter, transform, map_from_entries, col, struct

# 将Map转为键值对数组,过滤值不同的条目后转回Map
result_df = df.withColumn(
    "value_diff",
    map_from_entries(
        filter(
            transform(
                map_entries(col("sot")),
                lambda entry: struct(
                    entry.key.alias("key"),
                    entry.value.alias("old"),
                    col("metrics")[entry.key].alias("new")
                )
            ),
            lambda x: x.old != x.new
        )
    )
)

result_df.show(truncate=False)

逻辑拆解:

  1. map_entries(col("sot")):将sot的Map类型转为Array<Struct(key, value)>格式
  2. transform(...):遍历每个条目,关联metrics中对应key的值,生成包含key, old, new的结构体数组
  3. filter(...):筛选出old != new的差异条目
  4. map_from_entries(...):将过滤后的结构体数组转回Map类型

注意事项

  • 若sot和metrics是JSON字符串,需先用from_json函数转为Map类型再处理。
  • 内置函数方案性能优于UDF,Spark可对其做全链路优化,无需跨语言序列化。

内容的提问来源于stack exchange,提问作者Akshay Hazari

相关产品推荐
方舟 Agent Plan

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

最近更新时间:2026.07.18 22:12:28