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

PySpark如何通过字典列表匹配ID同时更新DataFrame多列值

方案1:修改UDF返回多值(适合小数据量场景)

你可以先把字典列表预处理成以id为key的映射字典,避免UDF每次调用都遍历整个列表,再让UDF返回结构体类型,一次带出key_1和key_2两个值:

# 先把列表转成字典,提升查找效率
id_map = {item['id']: (item['key_1'], item['key_2']) for item in data_list}

定义返回结构体的UDF:

from pyspark.sql.types import StructType, StructField, IntegerType
from pyspark.sql.functions import udf, col

# 定义UDF返回的结构体类型
return_cols_schema = StructType([
    StructField("key_1", IntegerType(), nullable=True),
    StructField("key_2", IntegerType(), nullable=True)
])

@udf(returnType=return_cols_schema)
def get_matched_cols(id_val):
    # 统一转字符串匹配,避免df的Id是数值类型和字典的字符串id不匹配
    return id_map.get(str(id_val), (None, None))

调用UDF替换两列:

df = df.withColumn("matched", get_matched_cols(col("Id"))) \
    .withColumn("Column_1", col("matched.key_1")) \
    .withColumn("Column_2", col("matched.key_2")) \
    .drop("matched")

# 查看结果
df.show()

方案2:使用Join实现(推荐,大数据量性能更好)

PySpark原生Join的性能远高于自定义UDF,更推荐把字典列表转成Spark DataFrame后关联更新:

# 把字典列表转成Spark DataFrame
map_df = spark.createDataFrame(data_list)

# 关联更新列,如果Id和id类型不匹配,可把a.Id转成字符串:col("a.Id").cast("string") == col("b.id")
df = df.alias("a") \
    .join(map_df.alias("b"), col("a.Id") == col("b.id"), how="left") \
    .select(
        col("a.Id"),
        # 如果未匹配需要保留原值,可把对应行替换为 coalesce(col("b.key_1"), col("a.Column_1")).alias("Column_1")
        col("b.key_1").alias("Column_1"),
        col("b.key_2").alias("Column_2")
    )

df.show()

内容的提问来源于stack exchange,提问作者Abdu Malikov

相关产品推荐
方舟 Agent Plan

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

最近更新时间:2026.10.06 08:30:03