如何将含出发地、目的地及JSON列的PySpark DataFrame转为嵌套Python字典?
解决方案
方法一:利用PySpark内置函数高效聚合(推荐大数据量场景)
先通过PySpark的分组聚合生成目标结构的DataFrame,再转换为Python字典:
from pyspark.sql import functions as F # 1. 生成"出发地-目的地"格式的复合路由键 df_with_route = df.withColumn( "route", F.concat(F.col("fs_origin"), F.lit("-"), F.col("fs_destination")) ) # 2. 按路由分组,将年月与对应JSON列表映射为键值对字典 aggregated_df = df_with_route.groupBy("route").agg( F.map_from_arrays( F.collect_list("year-month"), F.collect_list("JSON") ).alias("monthly_data") ) # 3. 转换为目标嵌套字典 result_dict = {row["route"]: row["monthly_data"] for row in aggregated_df.collect()}
补充:若JSON列为字符串类型
如果你的JSON列是未解析的字符串格式,需要先将其解析为数组类型:
from pyspark.sql.types import ArrayType, StructType, StructField, StringType # 根据实际JSON结构定义schema(示例仅包含fs_date字段) json_schema = ArrayType(StructType([ StructField("fs_date", StringType(), nullable=True) ])) # 解析JSON字符串为数组 df_parsed = df.withColumn("JSON", F.from_json(F.col("JSON"), json_schema)) # 后续步骤同方法一 df_with_route = df_parsed.withColumn( "route", F.concat(F.col("fs_origin"), F.lit("-"), F.col("fs_destination")) ) aggregated_df = df_with_route.groupBy("route").agg( F.map_from_arrays( F.collect_list("year-month"), F.collect_list("JSON") ).alias("monthly_data") ) result_dict = {row["route"]: row["monthly_data"] for row in aggregated_df.collect()}
方法二:Python原生遍历构建字典(适合小数据量场景)
将DataFrame数据拉取到Driver端后,通过遍历生成目标字典:
# 拉取所有数据到Driver端(数据量大时慎用) rows = df.collect() result_dict = {} for row in rows: route_key = f"{row['fs_origin']}-{row['fs_destination']}" month_key = row['year-month'] json_list = row['JSON'] # 初始化路由对应的内层字典 if route_key not in result_dict: result_dict[route_key] = {} # 赋值年月与对应JSON列表 result_dict[route_key][month_key] = json_list
内容的提问来源于stack exchange,提问作者Daniel Avigdor
相关产品推荐
相关产品推荐

