PySpark中如何将字典作为列值存储?合并ID关联DataFrame
解决Spark DataFrame将ID对应列转为字典列的问题
单个DataFrame处理方法
别用collect()把数据拉到Driver端,大数据场景下这会导致内存溢出,还没法高效关联回原DataFrame。直接用PySpark内置的create_map函数就能分布式完成需求:
- 导入所需函数:
from pyspark.sql import functions as F
- 提取非ID列,生成键值对映射:
# 假设你的DataFrame名为df,ID列名为"ID" non_id_cols = [col for col in df.columns if col != "ID"] # 构造create_map的参数:每个列名作为键,列值作为值 map_entries = [] for col_name in non_id_cols: map_entries.append(F.lit(col_name)) # 键:列名字符串 map_entries.append(F.col(col_name)) # 值:列对应的数据 # 添加字典列并保留ID列 result_df = df.withColumn("attrs", F.create_map(*map_entries)).select("ID", "attrs")
执行后,result_df的attrs列就是每个ID对应的其他列键值对组成的字典(Spark的Map类型,后续可通过cast("string")转成字符串格式,或直接用Map类型操作)。
两个DataFrame的处理与合并
因为两个DF列数不同,你可以分别处理每个DF生成字典列,再按ID关联:
# 处理第一个DF df1 non_id_cols1 = [col for col in df1.columns if col != "ID"] map_entries1 = [] for c in non_id_cols1: map_entries1.append(F.lit(c)) map_entries1.append(F.col(c)) df1_processed = df1.withColumn("attrs_df1", F.create_map(*map_entries1)).select("ID", "attrs_df1") # 处理第二个DF df2 non_id_cols2 = [col for col in df2.columns if col != "ID"] map_entries2 = [] for c in non_id_cols2: map_entries2.append(F.lit(c)) map_entries2.append(F.col(c)) df2_processed = df2.withColumn("attrs_df2", F.create_map(*map_entries2)).select("ID", "attrs_df2") # 按ID合并两个结果,outer join保留所有ID final_df = df1_processed.join(df2_processed, on="ID", how="outer")
为什么collect()的方法不可行
[row.asDict() for row in df.collect()]会把整个DataFrame的所有数据加载到Driver节点内存中,数据量稍大就会触发内存溢出。而且得到的字典列表无法直接和原DataFrame的ID做分布式关联,只能手动循环处理,效率极低,完全不适合大数据场景。
内容的提问来源于stack exchange,提问作者mohammed shoab
相关产品推荐
相关产品推荐

