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

DynamoDB导出S3的JSON加载到Spark时categories字段为null问题求助

问题根因

DynamoDB导出到S3的JSON文件自带数据类型标识前缀(M代表Map类型、N代表数值类型、S代表字符串类型),你当前定义的schema与实际JSON嵌套结构完全不匹配,因此Spark解析失败,categories字段返回null。

解决方案

步骤1:修正schema定义

按照实际JSON嵌套结构定义schema:

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

val schema = StructType(Array(
  // 适配DynamoDB导出的categories嵌套结构
  StructField("categories", StructType(Seq(
    StructField("M", MapType(StringType, StructType(Seq(
      StructField("N", DoubleType)
    ))))
  )), nullable = true),
  // 其他带类型标识的字段保留StringType即可,后续可按需提取实际值
  StructField("nextRequestAt", StringType, true),
  StructField("requestId", StringType, true),
  StructField("requestedAt", StringType, true),
  StructField("status", StringType, true),
  StructField("url", StringType, true),
  StructField("validUntil", StringType, true)
))

val df = spark.read.option("multiLine", true)
  .schema(schema)
  .json(s"${PATH}/*-load-dynamodb-data.json")

步骤2:提取扁平化的分类映射

如果你需要将categories转为key:分类名, value:数值的扁平Map,可通过以下方式转换:

import org.apache.spark.sql.functions.col

val resultDf = df.withColumn("categories_map", 
  // 提取M层的Map,再把每个value的N字段取出作为最终值
  col("categories.M").mapValues(_.getItem("N"))
)
// 可根据需要删除原始的嵌套categories字段
// .drop("categories")
可选优化

如果不想手动写schema,可以先让Spark自动推断结构,再基于推断结果做转换,适合数据结构不确定的场景:

// 读取小批量样本自动推断schema
val sampleDf = spark.read.option("multiLine", true).option("samplingRatio", 0.1).json(s"${PATH}/*-load-dynamodb-data.json")
// 打印推断的schema确认结构
sampleDf.printSchema()
// 后续基于打印的schema编写转换逻辑即可

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

相关产品推荐
方舟 Agent Plan

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

最近更新时间:2026.09.26 17:06:03