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
相关产品推荐
相关产品推荐

