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

如何使用PySpark将指定DataFrame转换为目标结构的字典?

PySpark实现DataFrame转指定结构字典的方法

步骤说明

  1. 过滤无效数据:剔除description字段为null的行,目标结果无需这类记录。
  2. 分组聚合生成嵌套字典:按description分组,将每组内的idx作为键、item_value作为值,生成嵌套的Map结构。
  3. 转换为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

相关产品推荐
方舟 Agent Plan

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

最近更新时间:2026.07.28 22:57:48