Spark DataFrame读取JSON时如何根据值动态设置key字段的数据类型
实现方案
步骤1:定义固定Schema读取JSON文件
由于原始JSON中key字段类型不固定,避免Spark自动推断Schema出现类型冲突,读取时先统一将key指定为字符串类型:
import org.apache.spark.sql.types._ // 定义values数组内部结构 val valueSchema = StructType(Seq( StructField("key", StringType, nullable = true), StructField("formId", IntegerType, nullable = true), StructField("occ", IntegerType, nullable = true), StructField("attachId", IntegerType, nullable = true) )) // 定义fields数组内部结构 val fieldSchema = StructType(Seq( StructField("fieldId", StringType, nullable = true), StructField("values", ArrayType(valueSchema), nullable = true) )) // 定义根节点Schema val rootSchema = StructType(Seq( StructField("id", StringType, nullable = true), StructField("name", StringType, nullable = true), StructField("fields", ArrayType(fieldSchema), nullable = true) )) // 读取JSON文件 val rawDF = spark.read.schema(rootSchema).json("你的JSON文件路径")
步骤2:实现类型转换逻辑
这里提供两种实现方案,可根据你的Spark版本选择:
方案A:兼容所有Spark版本(推荐,适配PIG读取)
自定义结构体存储转换后的key和对应类型标记,避免同列多类型导致的自动转型丢失信息:
import org.apache.spark.sql.Row import org.apache.spark.sql.functions.udf // 定义转换后的key结构体,包含类型标记和对应类型的值 val convertedKeyType = StructType(Seq( StructField("dataType", StringType, nullable = false), StructField("intValue", IntegerType, nullable = true), StructField("longValue", LongType, nullable = true), StructField("doubleValue", DoubleType, nullable = true), StructField("stringValue", StringType, nullable = true) )) // 定义类型转换UDF val keyTypeConvertUDF = udf((keyStr: String) => { if (keyStr == null) return null // 按精度从高到低尝试转换 try { val intVal = keyStr.toInt Row("int", intVal, null, null, null) } catch { case _: NumberFormatException => try { val longVal = keyStr.toLong Row("long", null, longVal, null, null) } catch { case _: NumberFormatException => try { val doubleVal = keyStr.toDouble Row("double", null, null, doubleVal, null) } catch { case _: NumberFormatException => Row("string", null, null, null, keyStr) } } } }, convertedKeyType)
方案B:Spark 3.4+ 简化方案
使用Spark内置的VariantType存储动态类型,无需自定义结构体:
import org.apache.spark.sql.functions._ // 直接用try_parse_json自动推断key的实际类型,存储为Variant val convertedDF = rawDF.withColumn("fields", transform(col("fields"), field => { field.withField("values", transform(field.getField("values"), valueItem => { valueItem.withField("key", try_parse_json(valueItem.getField("key"))) })) }))
步骤3:转换嵌套结构(仅方案A需要执行)
用transform函数处理嵌套的数组结构,替换key字段为转换后的结构体:
import org.apache.spark.sql.functions._ val convertedDF = rawDF.withColumn("fields", transform(col("fields"), field => { field.withField("values", transform(field.getField("values"), valueItem => { valueItem.withField("key", keyTypeConvertUDF(valueItem.getField("key"))) })) }))
步骤4:保存为Parquet文件
convertedDF.write.mode("overwrite").parquet("输出Parquet文件路径")
后续PIG脚本读取Parquet时,可通过dataType字段判断key的实际类型,读取对应的值即可。
内容的提问来源于stack exchange,提问作者amamagar
相关产品推荐
相关产品推荐

