如何在PySpark DataFrame中新增存储整行数据为字典的row_dict列
在PySpark DataFrame中新增整行转字典的列
方法:使用内置函数create_map实现
这种方法无需自定义UDF,性能更优,直接利用PySpark内置函数构造字典列:
- 导入所需模块
from pyspark.sql import functions as F from itertools import chain
- 构造字典列的表达式
假设你的DataFrame名为df,先获取所有列名,再通过create_map将列名(作为字典的key)和对应列值(作为字典的value)一一配对:
# 获取DataFrame的所有列名 all_columns = df.columns # 构造create_map的参数:交替传入列名字面量和列本身 map_expression = F.create_map(*chain(*[(F.lit(col_name), F.col(col_name)) for col_name in all_columns]))
- 新增
row_dict列
# 生成包含row_dict列的新DataFrame result_df = df.withColumn("row_dict", map_expression)
说明
- 生成的
row_dict列类型为MapType(StringType, AnyType),对应Python中的字典结构,key是列名字符串,value是对应行的原始值(支持所有PySpark数据类型,如字符串、数字、数组、结构体等) - 相比自定义UDF,这种方法完全基于PySpark内置优化,执行效率更高
内容的提问来源于stack exchange,提问作者pbh
相关产品推荐
相关产品推荐

