Spark-XML解析Windows-1251编码XML文件的异常问题咨询
问题分析与解决
为什么指定charset="cp1251"仍乱码?
在PySpark 2.4.0搭配的旧版spark-xml(约0.9.x版本)中,charset参数的处理存在逻辑冲突:
- Spark会先按照
charset指定的编码(这里是cp1251)把文件内容解码为Unicode字符串 - 后续底层的StAX XML解析器看到XML声明里的
encoding='WINDOWS-1251',会尝试将已解码的Unicode字符串再次当作cp1251字节流去解码,双重解码直接导致西里尔文乱码。
而将文件转成UTF-8后,即使指定charset="cp1251",Spark用cp1251读取UTF-8文件时刚好未触发双重解码冲突,才出现解析正常的巧合。
为什么Spark会忽略XML声明的编码?
Spark DataFrameReader的工作流程是先按指定/默认编码读取整个文件为文本流,再将解码后的内容交给spark-xml解析器处理。XML声明里的编码是给XML解析器用的,但此时字节流已经被Spark转换成Unicode字符串,解析器无法再基于原始字节流识别编码——因此XML声明的编码会被忽略,Spark读取阶段的编码优先级完全更高。
直接读取cp1251编码XML的可行方案
方案1:二进制读取后手动解析
绕开Spark的文本编码读取逻辑,直接读二进制内容交给XML解析器,让解析器自动识别XML声明的编码:
from pyspark.sql.functions import udf, explode from pyspark.sql.types import StructType, StructField, StringType, ArrayType import xml.etree.ElementTree as ET # 读取二进制文件 binary_df = spark.read.format("binaryFile").load("./cp1251") # 自定义UDF解析二进制内容(根据你的嵌套结构调整字段提取逻辑) def parse_xml_binary(content): root = ET.fromstring(content) product_list = [] for product in root.findall(".//product"): # 示例提取嵌套的name字段,可根据实际结构修改 name = product.find(".//name").text if product.find(".//name") is not None else None product_list.append({"product_name": name}) return product_list # 定义返回的StructType,匹配你的数据结构 product_schema = ArrayType(StructType([ StructField("product_name", StringType(), nullable=True) ])) parse_udf = udf(parse_xml_binary, product_schema) # 解析并展开结果 parsed_df = binary_df.select(explode(parse_udf("content")).alias("product")).select("product.*")
方案2:升级spark-xml版本
旧版spark-xml的编码处理逻辑存在缺陷,升级到0.12.x及以上版本后,修复了编码识别问题,支持正确联动XML声明的编码。升级后可直接读取:
# 无需指定charset,让解析器自动识别 df = spark.read.format('xml').options(rootTag='products', rowTag='product').load('./cp1251') # 或明确指定编码(新版本不会触发双重解码) df = spark.read.format('xml').options(rootTag='products', rowTag='product', charset="cp1251").load('./cp1251')
方案3:Spark内置API转码(应急方案)
如果无法升级依赖,可先用Spark将文件转码为UTF-8后再读取:
# 按cp1251读取文本,转码为UTF-8保存 text_df = spark.read.text("./cp1251", encoding="cp1251") text_df.write.text("./utf8_converted", encoding="utf-8") # 读取转码后的文件 df = spark.read.format('xml').options(rootTag='products', rowTag='product').load('./utf8_converted')
内容的提问来源于stack exchange,提问作者Владислав Черкасов
相关产品推荐
相关产品推荐

