Spark 2.x中解析表内XML列并拆分至单独列的方法
Spark 2.x解析XML列并拆分至单独字段的解决方案
1. 依赖准备
Spark 2.x本身没有内置XML解析能力,需要引入第三方库com.databricks:spark-xml,版本需与Spark、Scala版本匹配:
- Spark 2.x 搭配 Scala 2.11 时,使用
com.databricks:spark-xml_2.11:0.6.0 - 提交PySpark任务时,通过
--packages参数引入:pyspark --packages com.databricks:spark-xml_2.11:0.6.0 - 集群环境需确保该jar包已添加到Spark的classpath中。
2. 解析方案
假设你的DataFrame名为dfdata,存储XML字符串的列名为xml_col,以下分两种场景提供实现:
场景A:Spark 2.3及以上版本(推荐)
Spark 2.3引入了from_xml函数,可直接结合预定义Schema解析XML列,性能更优:
步骤1:定义XML对应的Schema
from pyspark.sql.types import StructType, StructField, StringType, IntegerType # 匹配XML结构的Schema xml_schema = StructType([ StructField("Items", StructType([ StructField("Item", StructType([ StructField("articleId", StringType()), StructField("inHandDifference", IntegerType()), StructField("articleName", StringType()), StructField("articleCode", StringType()), StructField("cappingLimit", IntegerType()), StructField("unitOfMeasure", StringType()), StructField("quantityInHand", IntegerType()), StructField("someid", StringType()), StructField("someotherid", StringType()) ])) ])) ])
步骤2:解析XML并提取字段
from pyspark.sql.functions import from_xml # 解析XML列得到嵌套结构 df_parsed = dfdata.withColumn("xml_parsed", from_xml(dfdata["xml_col"], xml_schema)) # 提取嵌套字段为单独列,并清理临时列 df_final = df_parsed.select( "*", df_parsed.xml_parsed.Items.Item.articleId.alias("article_id"), df_parsed.xml_parsed.Items.Item.inHandDifference.alias("in_hand_difference"), df_parsed.xml_parsed.Items.Item.articleName.alias("article_name"), df_parsed.xml_parsed.Items.Item.articleCode.alias("article_code"), df_parsed.xml_parsed.Items.Item.cappingLimit.alias("capping_limit"), df_parsed.xml_parsed.Items.Item.unitOfMeasure.alias("unit_of_measure"), df_parsed.xml_parsed.Items.Item.quantityInHand.alias("quantity_in_hand"), df_parsed.xml_parsed.Items.Item.someid.alias("some_id"), df_parsed.xml_parsed.Items.Item.someotherid.alias("some_other_id") ).drop("xml_col", "xml_parsed") # 查看解析结果 df_final.show()
场景B:Spark 2.3以下版本
若使用更早的Spark 2.x版本,可通过自定义UDF结合Python的xml.etree.ElementTree库解析:
步骤1:编写XML解析UDF
import xml.etree.ElementTree as ET from pyspark.sql.functions import udf from pyspark.sql.types import Row def parse_xml(xml_str): try: root = ET.fromstring(xml_str) item_node = root.find(".//Item") # 提取各字段值,注意类型转换 return Row( articleId=item_node.findtext("articleId"), inHandDifference=int(item_node.findtext("inHandDifference")) if item_node.findtext("inHandDifference") else None, articleName=item_node.findtext("articleName"), articleCode=item_node.findtext("articleCode"), cappingLimit=int(item_node.findtext("cappingLimit")) if item_node.findtext("cappingLimit") else None, unitOfMeasure=item_node.findtext("unitOfMeasure"), quantityInHand=int(item_node.findtext("quantityInHand")) if item_node.findtext("quantityInHand") else None, someid=item_node.findtext("someid"), someotherid=item_node.findtext("someotherid") ) except Exception: # 解析失败时返回空值,避免任务中断 return Row( articleId=None, inHandDifference=None, articleName=None, articleCode=None, cappingLimit=None, unitOfMeasure=None, quantityInHand=None, someid=None, someotherid=None ) # 绑定UDF的返回类型(复用之前定义的xml_schema中的Item结构) parse_xml_udf = udf(parse_xml, xml_schema["Items"]["Item"].dataType)
步骤2:应用UDF并提取字段
# 应用UDF解析XML列 df_parsed = dfdata.withColumn("xml_parsed", parse_xml_udf(dfdata["xml_col"])) # 提取字段并清理临时列 df_final = df_parsed.select( "*", df_parsed.xml_parsed.articleId.alias("article_id"), df_parsed.xml_parsed.inHandDifference.alias("in_hand_difference"), df_parsed.xml_parsed.articleName.alias("article_name"), df_parsed.xml_parsed.articleCode.alias("article_code"), df_parsed.xml_parsed.cappingLimit.alias("capping_limit"), df_parsed.xml_parsed.unitOfMeasure.alias("unit_of_measure"), df_parsed.xml_parsed.quantityInHand.alias("quantity_in_hand"), df_parsed.xml_parsed.someid.alias("some_id"), df_parsed.xml_parsed.someotherid.alias("some_other_id") ).drop("xml_col", "xml_parsed") df_final.show()
3. 关键注意事项
- 多Item节点处理:若XML中包含多个
<Item>节点,需将Schema中的Item字段改为ArrayType(StructType([...])),再用explode函数展开数组。 - 异常处理:实际业务中可能存在格式错误的XML,需在解析逻辑中添加异常捕获,避免任务失败。
- 版本匹配:确保spark-xml库的版本与Spark、Scala版本兼容,避免依赖冲突。
内容的提问来源于stack exchange,提问作者its_niks
相关产品推荐
相关产品推荐

