Spark DataFrame解析混合内容XML的字段提取及节点识别问题
解决Spark DataFrame解析XML时的Measure节点问题
看起来你遇到的两个问题都是因为Schema定义和Spark XML的配置没匹配上,我来帮你一步步解决:
问题根源分析
- 仅识别到一个Measure节点:你的Schema里把
Measure定义成了单个StructType,但XML里是多个<Measure>节点,应该用ArrayType来包裹,否则Spark只会读取第一个节点。 - 无法读取Answer/Question子字段且只能获取Measure文本:
<Measure>节点是混合内容(既有文本又有子节点),Spark XML默认不会自动处理这种情况,需要显式指定文本内容对应的字段名,同时在Schema里包含子节点字段。
修正后的解决方案
1. 调整自定义Schema
把Measure改成ArrayType,同时添加一个字段(比如text)来存储Measure的文本内容:
import org.apache.spark.sql.types._ def getCustomSchema(): StructType = { StructType(Array( StructField("QData", StructType(Array( StructField("Measure", ArrayType(StructType(Array( StructField("text", StringType, true), StructField("Answer", StringType, true), StructField("Question", StringType, true) ))), true) )), true) )) }
2. 修改XML读取代码
添加valueTag选项,告诉Spark把Measure节点的文本内容映射到我们定义的text字段:
val result = sc.read .format("com.databricks.spark.xml") .option("attributePrefix", "attr_") .option("valueTag", "text") // 指定文本内容对应的字段名 .schema(getCustomSchema) .load(filename.toString)
3. 修正映射逻辑(QDMapper)
现在Measure是数组类型,需要遍历每个Measure元素,而不是直接取单个Row:
case class QData(text: String, answer: String, question: String) case class QDMapper(){ def apply(row: Row): List[QData] = { val qDList = new scala.collection.mutable.ListBuffer[QData]() val qualData = row.getAs[Row]("QData") // 获取Measure数组 val measures = qualData.getAs[Seq[Row]]("Measure") measures.foreach(measureRow => { val text = measureRow.getAs[String]("text") val answer = measureRow.getAs[String]("Answer") val question = measureRow.getAs[String]("Question") qDList.append(QData(text, answer, question)) }) qDList.toList } }
4. 最终数据处理
val qDfTemp = result .mapPartitions(partition => { val mapper = new QDMapper() partition.flatMap(row => mapper(row)) }) .toDF()
这样调整后,你应该能读取到所有的Measure节点,并且每个节点的文本、Answer和Question字段都能正确获取了。
内容的提问来源于stack exchange,提问作者Eric Thomas
相关产品推荐
相关产品推荐

