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

如何将DataFrame中Map(key:Int,value:struct)列转为JSON对象

解决方案:将Spark DataFrame中的Map类型列转换为指定JSON格式

针对你需要把嵌套Map-Struct类型列转换为特定结构JSON的需求,我整理了Spark Scala和PySpark两种常用实现方案,完全匹配你的格式要求:

核心思路

先把Map列的每个键值对展开为独立行,再提取Struct中的字段并重命名为目标JSON的字段名,最后将每行数据转换为JSON字符串。注意你期望输出里VALUE_EN带双引号,我会在代码里加入这个处理,不需要的话可以直接移除。

1. Spark Scala 实现

先导入必要的Spark工具类:

import org.apache.spark.sql.functions._
import org.apache.spark.sql.types._

假设你的原DataFrame名为df,执行以下处理:

// 展开Map列,拆分出key和对应的Struct值
val processedDF = df
  .select(explode(col("PRODUCT_ID_FETR_MAP")).alias("map_entry"))
  .select(
    // 将Map的key转为字符串,对应JSON中的sortpriority字段
    col("map_entry.key").cast(StringType).alias("sortpriority"),
    // 严格映射Struct字段到目标JSON字段名
    col("map_entry.value.feat_id").alias("FEATURE_ID"),
    // 给VALUE_EN添加双引号,匹配你的期望格式
    concat(lit("\""), col("map_entry.value.feat_value_en"), lit("\"")).alias("VALUE_EN"),
    col("map_entry.value.feat_value_fr").alias("VALUE_FR"),
    col("map_entry.value.feat_def_key").alias("DEFN_KEY"),
    col("map_entry.value.feat_def_name_en").alias("DISPLAY_NAME_EN"),
    col("map_entry.value.feat_def_name_fr").alias("DISPLAY_NAME_FR"),
    col("map_entry.value.feat_def_sortpriority").alias("SORT_PRIORITY"),
    col("map_entry.value.feat_group_id").alias("feat_group_id"),
    col("map_entry.value.feat_grp_name_en").alias("feat_grp_name_en"),
    col("map_entry.value.feat_grp_name_fr").alias("feat_grp_name_fr")
  )

// 将每行数据转换为JSON字符串
val resultDF = processedDF.toJSON

2. PySpark 实现

先导入依赖函数:

from pyspark.sql import functions as F
from pyspark.sql.types import StringType

处理代码如下:

# 假设原DataFrame名为df
processed_df = df \
    .select(F.explode(F.col("PRODUCT_ID_FETR_MAP")).alias("map_entry")) \
    .select(
        F.col("map_entry.key").cast(StringType()).alias("sortpriority"),
        F.col("map_entry.value.feat_id").alias("FEATURE_ID"),
        # 给VALUE_EN添加双引号
        F.concat(F.lit('"'), F.col("map_entry.value.feat_value_en"), F.lit('"')).alias("VALUE_EN"),
        F.col("map_entry.value.feat_value_fr").alias("VALUE_FR"),
        F.col("map_entry.value.feat_def_key").alias("DEFN_KEY"),
        F.col("map_entry.value.feat_def_name_en").alias("DISPLAY_NAME_EN"),
        F.col("map_entry.value.feat_def_name_fr").alias("DISPLAY_NAME_FR"),
        F.col("map_entry.value.feat_def_sortpriority").alias("SORT_PRIORITY"),
        F.col("map_entry.value.feat_group_id").alias("feat_group_id"),
        F.col("map_entry.value.feat_grp_name_en").alias("feat_grp_name_en"),
        F.col("map_entry.value.feat_grp_name_fr").alias("feat_grp_name_fr")
    )

# 转换为JSON格式的结果DataFrame
result_df = processed_df.toJSON()

补充提示

  • 如果原DataFrame有其他需要保留的列(比如产品ID),只需在第一步select时添加对应列即可,例如:df.select(col("PRODUCT_ID"), explode(col("PRODUCT_ID_FETR_MAP")).alias("map_entry"))
  • 若不需要VALUE_EN带双引号,直接替换为F.col("map_entry.value.feat_value_en").alias("VALUE_EN")即可
  • 执行后,resultDF(或result_df)的每一行就是你需要的JSON字符串格式

内容的提问来源于stack exchange,提问作者Chandra

相关产品推荐
方舟 Agent Plan

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

最近更新时间:2026.05.27 07:06:09