如何使用PySpark将指定DataFrame转换为目标结构的字典?
PySpark实现DataFrame转指定结构字典的方法
步骤说明
- 过滤无效数据:剔除
description字段为null的行,目标结果无需这类记录。 - 分组聚合生成嵌套字典:按
description分组,将每组内的idx作为键、item_value作为值,生成嵌套的Map结构。 - 转换为Python字典:将聚合后的Spark DataFrame结果转换为Python原生字典。
完整代码实现
from pyspark.sql import SparkSession from pyspark.sql.functions import col, map_from_entries, collect_list, struct # 初始化SparkSession spark = SparkSession.builder.appName("DataFrameToDict").getOrCreate() # 构造示例DataFrame data = [ ("A", 0.25, "2023-03-01T17:20:00.000+0000", 0, "unitA&"), ("B", 0.34, "2023-03-01T17:20:00.000+0000", 0, None), ("C", 0.34, "2023-03-01T17:20:00.000+0000", 0, "C&"), ("D", 0.45, "2023-03-01T17:20:00.000+0000", 0, "unitD&"), ("A", 0.3, "2023-03-01T17:25:00.000+0000", 1, "unitA&"), ("B", 0.54, "2023-03-01T17:25:00.000+0000", 1, None), ("C", 0.3, "2023-03-01T17:25:00.000+0000", 1, "C&"), ("D", 0.21, "2023-03-01T17:25:00.000+0000", 1, "unitD& "), ("A", 0.54, "2023-03-01T17:30:00.000+0000", 2, "unitA& ") ] df = spark.createDataFrame(data, ["item_name", "item_value", "timestamp", "idx", "description"]) # 1. 过滤null的description行 filtered_df = df.filter(col("description").isNotNull()) # 2. 分组聚合生成嵌套Map aggregated_df = filtered_df.groupBy("description") \ .agg(map_from_entries(collect_list(struct(col("idx"), col("item_value")))).alias("value_map")) # 3. 转换为Python字典 result_dict = {row.description: row.value_map for row in aggregated_df.collect()} print(result_dict)
代码解释
- 过滤步骤:使用
filter(col("description").isNotNull())剔除description为空的记录,对应原数据中item_name为B的行。 - 聚合步骤:
struct(col("idx"), col("item_value")):将idx和item_value组合成键值对结构体。collect_list(...):收集每个description分组内的所有键值对结构体,形成列表。map_from_entries(...):将键值对列表转换为Spark的Map类型,实现{idx: item_value}的嵌套结构。
- 转字典:通过
collect()获取所有行数据,再用字典推导式将每行的description作为键、value_map作为值,生成最终的Python字典。
输出结果
运行代码后,输出将与目标结构一致(若需合并末尾带空格的相同description,可对字段做trim()处理后再分组):
{'unitA&': {0: 0.25, 1: 0.3}, 'C&': {0: 0.34, 1: 0.3}, 'unitD&': {0: 0.45}, 'unitD& ': {1: 0.21}, 'unitA& ': {2: 0.54}}
内容的提问来源于stack exchange,提问作者MMV
相关产品推荐
相关产品推荐

