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

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的结构与你需要的表结构完全一致:

namesurnameadress.cityadress.street
bbde
bcde
bcfg

补充说明

  • 批处理场景只需将readStream替换为read,调整数据读取方式(比如从文件读取JSON)即可。
  • 用typeof判断字段类型比is_array更准确,因为单个地址是struct类型,数组是array类型。
  • 不要一开始就用强Schema读取,否则Spark会将类型不匹配的行标记为null或直接过滤,造成数据丢失。

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

相关产品推荐
方舟 Agent Plan

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

最近更新时间:2026.07.02 21:17:11