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

在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

相关产品推荐
方舟 Agent Plan

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

最近更新时间:2026.06.23 10:18:26