如何使用Spark与Java将XML解析转换为扁平化单行数据表
实现方案
问题原因
你原来的实现存在3个核心问题:
- 递归调用flattenSchema返回的列没有合并到当前结果集,递归生成的列全部丢失
- 没有处理
ArrayType类型字段:XML中的重复标签(比如、 、 )会被spark-xml解析为数组类型,你原逻辑只处理了StructType - 未处理XML属性到列的映射:spark-xml默认会给标签属性加
_前缀,同时重复标签的属性和值需要先做行转列才能匹配你期望的输出结构
步骤1:修正XML读取配置
首先调整读取参数,确保属性和标签值正确解析:
import static org.apache.spark.sql.functions.*; Dataset<Row> rawDf = spark.read() .format("com.databricks.spark.xml") .option("rootTag", "DMessage") .option("rowTag", "Message") // 每条Message对应一行数据 .option("attributePrefix", "_") // 可自定义属性前缀,默认即为_ .load(filePath.toString());
读取后的数据为嵌套的Struct+Array混合结构。
步骤2:预处理重复数组节点,转成键值对结构
针对你XML中三类重复标签,先做行转列处理,把类型属性作为键、标签内容作为值:
// 处理Identifiers下的Id数组 Dataset<Row> step1 = rawDf.withColumn("IdentifierMap", map_from_entries(expr("transform(IP.Identifiers.Id, x -> struct(x._Type as key, x._VALUE as value))")) ); // 处理Name下的NameComponent数组 Dataset<Row> step2 = step1.withColumn("NameMap", map_from_entries(expr("transform(IP.Names.Name[0].NameComponent, x -> struct(x._NameComponentType as key, x._VALUE as value))")) ); // 处理Parents下的Parent数组 Dataset<Row> step3 = step2.withColumn("ParentList", expr("IP.Parents.Parent")) .withColumn("Parent_T1_Map", map_from_entries(expr("transform(filter(ParentList, x -> x._ParentType = 'T1')[0].Id, x -> struct(x._Type as key, x._VALUE as value))")) ) .withColumn("Parent_T2_Map", map_from_entries(expr("transform(filter(ParentList, x -> x._ParentType = 'T2')[0].Id, x -> struct(x._Type as key, x._VALUE as value))")) );
步骤3:通用扁平化方法(支持Struct+Array混合结构)
修正后的递归方法,同时处理StructType和ArrayType:
private static List<Column> flattenSchema(StructType schema, String prefix) { List<Column> out = new ArrayList<>(); Arrays.stream(schema.fields()).forEach(f -> { String colName = prefix == null ? f.name() : prefix + "." + f.name(); if (f.dataType() instanceof StructType) { // 合并递归返回的列到当前结果集 out.addAll(flattenSchema((StructType) f.dataType(), colName)); } else if (f.dataType() instanceof ArrayType) { // 简单数组可直接转字符串,复杂数组建议提前做行转列处理 out.add(col(colName).cast("string").alias(colName.replace(".", "_"))); } else { out.add(col(colName).alias(colName.replace(".", "_"))); } }); return out; }
步骤4:生成最终目标表
你可以手动选取匹配期望结构的列,也可以调用通用方法全量扁平化:
// 手动选列,完全匹配你期望的输出结构(注:你示例输出中的Parent_T1_ID2重复为笔误,已修正为Parent_T2_D2) Dataset<Row> finalDf = step3.select( col("Header._Action").alias("Action"), col("Header._UpdatedBy").alias("UpdatedBy"), col("IdentifierMap.ID").alias("Identifier_ID"), col("IdentifierMap.ID2").alias("Identifier_ID2"), col("NameMap.FullName").alias("FullName"), col("NameMap.ShortName").alias("ShortName"), col("Parent_T1_Map.ID").alias("Parent_T1_ID"), col("Parent_T1_Map.ID2").alias("Parent_T1_ID2"), col("Parent_T2_Map.ID").alias("Parent_T2_ID"), col("Parent_T2_Map.D2").alias("Parent_T2_D2") ); // 全量自动扁平化可直接调用: // List<Column> flattenCols = flattenSchema(step3.schema(), null); // Dataset<Row> finalDf = step3.select(flattenCols);
执行finalDf.show()即可得到你需要的单行数据表。
内容的提问来源于stack exchange,提问作者babboon
相关产品推荐
相关产品推荐

