如何将PySpark DataFrame按fs_destination分区转为JSON列表字典
解决方案:PySpark DataFrame转换为按目的地分区的JSON字典列表
步骤说明
- 将每行数据转为JSON字符串:使用PySpark的
to_json函数,把DataFrame的每一行转换成标准JSON格式的字符串。 - 按目的地分组并收集列表:以
fs_destination为分组键,将每个分组下的所有JSON字符串收集成列表。 - 转换为Python字典:将分组后的结果从DataFrame转为Python字典,得到目标格式。
完整代码示例
from pyspark.sql import functions as F # 假设你的DataFrame名为df # 步骤1:添加JSON字符串列 df_with_json = df.withColumn("json_str", F.to_json(F.struct(df.columns))) # 步骤2:分组收集JSON列表 grouped_df = df_with_json.groupBy("fs_destination").agg(F.collect_list("json_str").alias("json_list")) # 步骤3:转换为Python字典 result_dict = {row.fs_destination: row.json_list for row in grouped_df.collect()} # 打印结果 print(result_dict)
输出结果
运行代码后,会得到符合要求的字典格式:
{ "AUH": [ '{"fs_date":"2022-06-01T00:00:00","ss_date":"2022-06-02T00:00:00","fs_origin":"TLV","fs_destination":"AUH","price":681.0715}', '{"fs_date":"2022-06-01T00:00:00","ss_date":"2022-06-03T00:00:00","fs_origin":"TLV","fs_destination":"AUH","price":406.46}' ], "BOM": [ '{"fs_date":"2022-06-01T00:00:00","ss_date":"2022-06-02T00:00:00","fs_origin":"TLV","fs_destination":"BOM","price":545.7715}', '{"fs_date":"2022-06-01T00:00:00","ss_date":"2022-06-03T00:00:00","fs_origin":"TLV","fs_destination":"BOM","price":372.435}' ] }
内容的提问来源于stack exchange,提问作者Daniel Avigdor
相关产品推荐
相关产品推荐

