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

PySpark读取含UUID的Parquet文件异常问题求助

解决Azure Synapse中PySpark读取含UUID的Parquet文件问题

问题根源

Parquet文件中的UUID字段通常以FIXED_LEN_BYTE_ARRAY(16)格式存储,PySpark默认无法直接将该类型映射为UUID字符串,导致不指定Schema时抛出Illegal Parquet type: FIXED_LEN_BYTE_ARRAY错误,指定Schema为StringType时则返回乱码。

解决方案

方案一:指定Schema并通过UDF转换二进制为UUID

先将UUID字段定义为BinaryType读取,再用Python的uuid模块将二进制数据转为标准UUID字符串:

from pyspark.sql.types import StructType, StructField, StringType, TimestampType, IntegerType, BinaryType
import uuid
from pyspark.sql.functions import udf

SourceFilePath = f'abfss://datalake@{AzureBlobStorageAccountName}.dfs.core.windows.net/Bronze/ServiceBus/{SourceFileName}'

# 定义包含BinaryType的Schema
schema = StructType([
    StructField('Created', TimestampType(), True),
    StructField('Entity', StringType(), True),
    StructField('EntityId', IntegerType(), True),
    StructField('Version', IntegerType(), True),
    StructField('CorrelationId', BinaryType(), True)
])

# 读取Parquet文件
dfBronze = spark.read.format("parquet").schema(schema).load(SourceFilePath)

# 定义UDF转换二进制为UUID字符串
binary_to_uuid = udf(lambda x: str(uuid.UUID(bytes=x)) if x is not None else None, StringType())

# 转换CorrelationId字段
dfBronze = dfBronze.withColumn('CorrelationId', binary_to_uuid('CorrelationId'))

dfBronze.createOrReplaceTempView("Bronze")
dfBronze.show()

方案二:使用Spark内置函数(Spark 3.0+)

如果你的Spark版本是3.0及以上,可以用内置的uuid_from_bytes函数,无需自定义UDF,性能更优:

from pyspark.sql.types import StructType, StructField, StringType, TimestampType, IntegerType, BinaryType
from pyspark.sql.functions import uuid_from_bytes

SourceFilePath = f'abfss://datalake@{AzureBlobStorageAccountName}.dfs.core.windows.net/Bronze/ServiceBus/{SourceFileName}'

schema = StructType([
    StructField('Created', TimestampType(), True),
    StructField('Entity', StringType(), True),
    StructField('EntityId', IntegerType(), True),
    StructField('Version', IntegerType(), True),
    StructField('CorrelationId', BinaryType(), True)
])

dfBronze = spark.read.format("parquet").schema(schema).load(SourceFilePath)

# 用内置函数转换二进制为UUID
dfBronze = dfBronze.withColumn('CorrelationId', uuid_from_bytes('CorrelationId'))

dfBronze.createOrReplaceTempView("Bronze")
dfBronze.show()

备选方案:关闭向量化读取(不推荐)

如果不想提前定义Schema,可以通过关闭Parquet向量化读取让Spark自动处理该类型,但会影响读取性能:

spark.conf.set("spark.sql.parquet.enableVectorizedReader", "false")

SourceFilePath = f'abfss://datalake@{AzureBlobStorageAccountName}.dfs.core.windows.net/Bronze/ServiceBus/{SourceFileName}'
dfBronze = spark.read.load(SourceFilePath, format='parquet').cache()

dfBronze.createOrReplaceTempView("Bronze")

内容的提问来源于stack exchange,提问作者user2197446

相关产品推荐
方舟 Agent Plan

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

最近更新时间:2026.07.13 09:23:19