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
相关产品推荐
相关产品推荐

