在Databricks的PySpark中解压SQL Server导出Parquet的压缩列
在Databricks PySpark中解压SQL Server导出的Parquet压缩列
场景背景
- SQL Server侧:原XML字段通过
COMPRESS(CAST([unzipped] AS VARCHAR(MAX)))压缩为varbinary(max)类型,采用gzip算法,可通过CAST(DECOMPRESS([zipped]) AS XML)还原。 - 导出环节:通过Azure Data Factory复制活动将表导出为Parquet文件时,
varbinary类型映射为Parquet的BINARY类型,存储在数据湖后由Databricks访问。 - 现状:Databricks读取后的DataFrame中,目标列类型为
StructField('column_name', BinaryType(), True),需要将其解压为String或XML类型。
解决方案
通过Python的gzip模块编写自定义UDF(用户定义函数),对二进制列进行解压,并将解压后的字节转为UTF-8字符串(这是解决显示异常的关键步骤)。
1. 导入依赖模块
import gzip from io import BytesIO from pyspark.sql.functions import udf from pyspark.sql.types import StringType
2. 编写解压UDF
def decompress_gzip_binary(binary_data): # 处理空值,避免空指针异常 if binary_data is None: return None try: # 用BytesIO包装二进制数据,通过gzip解压 with gzip.GzipFile(fileobj=BytesIO(binary_data), mode='rb') as f: decompressed_bytes = f.read() # 将字节转为UTF-8字符串(匹配SQL Server中CAST为VARCHAR(MAX)的编码) return decompressed_bytes.decode("UTF-8") except Exception as e: # 捕获解压异常,避免任务中断 print(f"解压列失败: {str(e)}") return None # 注册UDF,指定返回类型为StringType decompress_udf = udf(decompress_gzip_binary, StringType())
3. 应用UDF到DataFrame
假设读取后的DataFrame为df,压缩列名为zipped_column,执行以下代码生成解压后的列:
df_decompressed = df.withColumn("unzipped_xml_str", decompress_udf(df["zipped_column"]))
4. 可选:转为XML结构化对象(按需处理)
如果需要将字符串进一步解析为XML对象,可以额外编写UDF处理:
import xml.etree.ElementTree as ET def parse_xml_string(xml_str): if xml_str is None: return None try: # 解析XML字符串为ElementTree对象 root = ET.fromstring(xml_str) # 可根据需求返回结构化数据或保留原字符串,示例中返回原字符串 return xml_str except Exception as e: print(f"解析XML失败: {str(e)}") return None parse_xml_udf = udf(parse_xml_string, StringType()) df_with_xml = df_decompressed.withColumn("unzipped_xml", parse_xml_udf(df_decompressed["unzipped_xml_str"]))
关键注意事项
- 必须执行UTF-8解码:SQL Server是将XML转为
VARCHAR(MAX)(UTF-8编码)后再压缩,因此解压后的字节必须通过.decode("UTF-8")转为字符串,否则会出现乱码或显示异常。 - 空值与异常处理:加入空值判断和异常捕获,避免因个别损坏数据导致整个Spark任务失败。
内容的提问来源于stack exchange,提问作者Gam
相关产品推荐
相关产品推荐

