如何将航班数据PySpark DataFrame转换为指定嵌套字典?
将PySpark航班DataFrame转换为嵌套字典
实现步骤
直接通过PySpark的分组与Map转换逻辑即可实现需求,具体操作如下:
- 生成组合键:把出发地+目的地拼接成航班路线键(用
-分隔),出发日期+返程日期拼接成日期对键(用:分隔) - 分组构建Map:按航班路线分组,将每组内的日期对与价格映射为键值对Map
- 转成Python字典:把分组结果拉取到Driver端,用字典推导式生成目标嵌套结构
代码实现
假设你的DataFrame名为flight_df,执行以下代码:
from pyspark.sql import functions as F # 生成组合键列 transformed_df = flight_df.withColumn( "route", F.concat(F.col("fs_origin"), F.lit("-"), F.col("fs_destination")) ).withColumn( "date_pair", F.concat(F.col("fs_date"), F.lit(":"), F.col("ss_date")) ).select("route", "date_pair", "price") # 分组并将每组数据转为Map grouped_df = transformed_df.groupBy("route").agg( F.map_from_entries(F.collect_list(F.struct("date_pair", "price"))).alias("price_map") ) # 转换为目标嵌套字典 result_dict = {row.route: row.price_map for row in grouped_df.collect()}
代码说明
concat函数用于拼接字符串生成自定义键map_from_entries配合collect_list+struct,将分组内的日期对和价格批量转为Map结构collect()将分布式数据拉取到本地后,通过字典推导式直接生成所需的嵌套字典
内容的提问来源于stack exchange,提问作者Daniel Avigdor
相关产品推荐
相关产品推荐

