PySpark解析多XML文件仅提取部分数据的问题求助
Spark XML解析问题:仅提取到单条数据的原因与解决方法
环境与场景
- 技术栈:Spark 2.3.2 + Python 3.7
- 依赖:使用
com.databricks:spark-xml_2.11:0.7.0解析XML文件 - 问题现象:XML包含2个
<xocs:doc>节点,初始读取能获取2条数据,但提取ref-info及对应eid时,仅能得到eid=85082880163的数据;单独处理eid=85082880158的XML时可正常提取
初始读取代码(正常工作)
os.environ['PYSPARK_SUBMIT_ARGS'] = '--packages com.databricks:spark-xml_2.11:0.7.0 pyspark-shell' conf = pyspark.SparkConf() sc = SparkSession.builder.config(conf=conf).getOrCreate() spark = SQLContext(sc) dfSample = (spark.read.format("xml").option("rowTag", "xocs:doc") .load(r"sample.xml"))
提取数据的问题代码
(dfSample. withColumn("metaExp", F.explode(F.array("xocs:meta"))). withColumn("eid", F.col("metaExp.xocs:eid")). select("eid","xocs:item"). withColumn("xocs:itemExp", F.explode(F.array("xocs:item"))). withColumn("item", F.col("xocs:itemExp.item")). withColumn("itemExp", F.explode(F.array("item"))). withColumn("bibrecord", F.col("item.bibrecord")). withColumn("bibrecordExp", F.explode(F.array("bibrecord"))). withColumn("tail", F.col("bibrecord.tail")). withColumn("tailExp", F.explode(F.array("tail"))). withColumn("bibliography", F.col("tail.bibliography")). withColumn("bibliographyExp", F.explode(F.array("bibliography"))). withColumn("reference", F.col("bibliography.reference")). withColumn("referenceExp", F.explode(F.array("reference"))). withColumn("ref-infoExp", F.explode(F.col("reference.ref-info"))). withColumn("authors", F.explode(F.col("ref-infoExp.ref-authors.author"))). withColumn("py", (F.col("ref-infoExp.ref-publicationyear._first"))). withColumn("so", (F.col("ref-infoExp.ref-sourcetitle"))). withColumn("ti", (F.col("ref-infoExp.ref-title"))). drop("xocs:item", "xocs:itemExp", "item", "itemExp", "bibrecord", "bibrecordExp", "tail", "tailExp", "bibliography", "bibliographyExp", "reference", "referenceExp").show())
问题原因
核心问题出在嵌套层级的数组处理逻辑:
- 强制用
F.array()包裹字段后再explode,会导致单元素字段与数组字段的处理逻辑冲突,部分记录的字段结构不匹配时会被隐性过滤 - 多轮
explode过程中,若某条记录的某个嵌套层级无对应字段(或字段为空),explode会直接丢弃这条记录——比如eid=85082880158的记录可能在某个层级的字段结构与另一条不同,导致被过滤
解决方法
1. 用explode_outer()替代explode(),避免过滤数据
explode_outer()会保留原记录,即使字段为空或非数组,不会丢失数据。修改后的核心逻辑如下:
from pyspark.sql import functions as F (dfSample. # 处理meta层级,保留所有记录 withColumn("metaExp", F.explode_outer("xocs:meta")). withColumn("eid", F.col("metaExp.xocs:eid")). select("eid","xocs:item"). # 逐层用explode_outer替代explode,避免过滤无对应字段的记录 withColumn("xocs:itemExp", F.explode_outer("xocs:item")). withColumn("item", F.col("xocs:itemExp.item")). withColumn("itemExp", F.explode_outer("item")). withColumn("bibrecord", F.col("itemExp.bibrecord")). withColumn("bibrecordExp", F.explode_outer("bibrecord")). withColumn("tail", F.col("bibrecordExp.tail")). withColumn("tailExp", F.explode_outer("tail")). withColumn("bibliography", F.col("tailExp.bibliography")). withColumn("bibliographyExp", F.explode_outer("bibliography")). withColumn("reference", F.col("bibliographyExp.reference")). withColumn("referenceExp", F.explode_outer("reference")). withColumn("ref-infoExp", F.explode_outer(F.col("referenceExp.ref-info"))). # 提取字段时处理可能的空值 withColumn("authors", F.explode_outer(F.col("ref-infoExp.ref-authors.author"))). withColumn("py", F.col("ref-infoExp.ref-publicationyear._first")). withColumn("so", F.col("ref-infoExp.ref-sourcetitle")). withColumn("ti", F.col("ref-infoExp.ref-title")). # 清理冗余字段 drop("xocs:item", "xocs:itemExp", "item", "itemExp", "bibrecord", "bibrecordExp", "tail", "tailExp", "bibliography", "bibliographyExp", "reference", "referenceExp").show())
2. 提前校验字段结构差异
处理前先检查不同记录的字段结构是否一致,针对性调整逻辑:
# 查看每条记录的xocs:item字段是否为数组类型 dfSample.select( "eid", F.array_contains(F.schema_of_json(F.to_json("xocs:item")), "array").alias("is_array") ).show()
若存在单元素与数组字段的差异,可通过F.when判断后再决定是否用F.array包裹字段。
3. 批量处理建议
针对数千个XML文件的场景:
- 直接读取整个目录(
spark.read.xml().load("xml_directory/")),无需手动合并文件 - 抽样查看不同文件的字段结构,确保逻辑覆盖所有可能的结构差异
- 全程优先使用
explode_outer(),避免隐性数据丢失
内容的提问来源于stack exchange,提问作者mlee_jordan
相关产品推荐
相关产品推荐

