Spark解析Kafka中XML转JSON单/多地址数据异常求助
解决Spark解析单个/数组混合JSON字段的问题
你的问题核心是JSON中adress字段同时存在单个对象和数组两种格式,而Spark的强类型Schema无法兼容这两种情况。下面是具体的解决步骤:
1. 先修正Schema的错误
你原代码中的Schema定义了"a"字段,但实际JSON里的地址字段是"adress",而且直接指定ArrayType会导致单个对象的JSON解析失败(类型不匹配)。所以第一步要避免一开始就用强Schema约束,先以弱类型读取数据。
2. 统一adress字段为数组格式
读取数据后,通过Spark内置函数判断adress的类型,将单个对象包装为数组,确保所有行的adress都是数组类型:
import org.apache.spark.sql.functions._ import org.apache.spark.sql.types._ // 假设从Kafka读取数据,先解析JSON字符串(批处理场景替换为对应read方式即可) val rawDf = spark.readStream .format("kafka") .option("kafka.bootstrap.servers", "your-kafka-server:9092") .option("subscribe", "your-topic") .load() .selectExpr("CAST(value AS STRING) AS json_str") // 自动推断JSON结构并解析 .select(from_json(col("json_str"), schema_of_json(col("json_str"))).alias("data")) .select("data.*") // 统一adress格式:单个struct对象包装为数组,数组则保持原样 val unifiedDf = rawDf.withColumn( "adress", when(typeof(col("adress")) === "struct", array(col("adress"))) .otherwise(col("adress")) )
3. 展开数组并提取目标字段
使用explode函数将数组中的每个地址对象拆分为单独的行,再提取所需字段对齐目标表结构:
// 展开adress数组,每个地址对象对应一行 val explodedDf = unifiedDf.select( col("name"), col("surname"), explode(col("adress")).alias("adress_obj") ) // 提取地址字段,生成目标表结构 val finalDf = explodedDf.select( col("name"), col("surname"), col("adress_obj.city").alias("adress.city"), col("adress_obj.street").alias("adress.street") )
4. 验证结果
执行后finalDf的结构与你需要的表结构完全一致:
| name | surname | adress.city | adress.street |
|---|---|---|---|
| b | b | d | e |
| b | c | d | e |
| b | c | f | g |
补充说明
- 批处理场景只需将
readStream替换为read,调整数据读取方式(比如从文件读取JSON)即可。 - 用
typeof判断字段类型比is_array更准确,因为单个地址是struct类型,数组是array类型。 - 不要一开始就用强Schema读取,否则Spark会将类型不匹配的行标记为null或直接过滤,造成数据丢失。
内容的提问来源于stack exchange,提问作者Eupaphios
相关产品推荐
相关产品推荐

