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

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

相关产品推荐
方舟 Agent Plan

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

最近更新时间:2026.06.15 05:54:51