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

Azure Data Factory读取Parquet文件遇非原始类型问题求助

ADF读取Databricks写入的Parquet文件时非原始类型报错问题

在Azure Data Factory(ADF)中读取由Databricks Notebook写入Azure存储的Parquet文件时,加载数据集和执行复制活动均触发非原始类型报错。以下是Databricks侧的输入数据格式、写入代码及问题场景:

Databricks侧输入数据与代码

# 输入数据格式
[
    {'Details': {'Input': {'id': '1', 'name': 'asdsdasd', 'a1': None, 'a2': None, 'c': None, 's': None, 'c1': None, 'z': None}, 'Output': '{"msg":"some error"}'}, 'Failure': '{"msg":"error"}', 's': 'f'},
    {'Details': {'Input': {'id': '2', 'name': 'sadsadsad', 'a1': 'adsadsad', 'a2': 'sssssss', 'c': 'cccc', 's': 'test', 'c1': 'ind', 'z': '22222'}, 'Output': '{"s":"2"}'}, 'Failure': '', 's': 's'}
]

# 写入Parquet的代码
from pyspark.sql.functions import to_json,col
from pyspark.sql.types import StructType, StructField, StringType

schema = StructType([
    StructField("id", StringType(), nullable=False),
    StructField("name", StringType(), nullable=True),
    StructField("desc", StringType(), nullable=True),
    StructField("Details", StructType([
        StructField("Input", StringType(), nullable=True),
        StructField("Output", StringType(), nullable=True)
    ])),
    StructField("Failure", StringType(), nullable=True),
    StructField("s", StringType(), nullable=True)
])

newJson = [
    {'Details': {'Input': {'id': '1', 'name': 'asdsdasd', 'a1': None, 'a2': None, 'c': None, 's': None, 'c1': None, 'z': None}, 'Output': '{"msg":"some error"}'}, 'Failure': '{"msg":"error"}', 's': 'f'},
    {'Details': {'Input': {'id': '2', 'name': 'sadsadsad', 'a1': 'adsadsad', 'a2': 'sssssss', 'c': 'cccc', 's': 'test', 'c1': 'ind', 'z': '22222'}, 'Output': '{"s":"2"}'}, 'Failure': '', 's': 's'}
]

df=spark.createDataFrame(data=newJson,schema=schema)
display(df)
df.coalesce(1).write.parquet(f"{adls_url}/test/", mode="overwrite")

问题根源

核心矛盾在于Schema定义与实际数据类型不匹配:

  • 代码中定义Details.Input为StringType,但实际传入的是嵌套字典(Struct类型数据)
  • Spark在写入Parquet时,会优先根据实际数据推断类型并覆盖Schema定义,导致最终Parquet文件中Details.Input是嵌套Struct类型,而非预期的字符串类型
  • ADF对Parquet的嵌套非原始类型处理存在兼容性限制,无法直接识别这种嵌套结构,因此触发报错

解决方案

方案1:修正Databricks代码,确保Schema与数据类型一致

将嵌套的Input字典转换为JSON字符串,严格匹配Schema中定义的StringType:

from pyspark.sql.functions import to_json, col, struct
from pyspark.sql.types import StructType, StructField, StringType

schema = StructType([
    StructField("id", StringType(), nullable=False),
    StructField("name", StringType(), nullable=True),
    StructField("desc", StringType(), nullable=True),
    StructField("Details", StructType([
        StructField("Input", StringType(), nullable=True),
        StructField("Output", StringType(), nullable=True)
    ])),
    StructField("Failure", StringType(), nullable=True),
    StructField("s", StringType(), nullable=True)
])

newJson = [
    {'Details': {'Input': {'id': '1', 'name': 'asdsdasd', 'a1': None, 'a2': None, 'c': None, 's': None, 'c1': None, 'z': None}, 'Output': '{"msg":"some error"}'}, 'Failure': '{"msg":"error"}', 's': 'f'},
    {'Details': {'Input': {'id': '2', 'name': 'sadsadsad', 'a1': 'adsadsad', 'a2': 'sssssss', 'c': 'cccc', 's': 'test', 'c1': 'ind', 'z': '22222'}, 'Output': '{"s":"2"}'}, 'Failure': '', 's': 's'}
]

# 创建DataFrame后,将Input字段转换为JSON字符串
df = spark.createDataFrame(data=newJson)
df = df.withColumn("Details", 
                   struct(
                       to_json(col("Details.Input")).alias("Input"),
                       col("Details.Output")
                   ))

# 强制匹配Schema,确保类型严格一致
df = df.select("id", "name", "desc", "Details", "Failure", "s").cast(schema)

display(df)
df.coalesce(1).write.parquet(f"{adls_url}/test/", mode="overwrite")

方案2:在ADF中适配嵌套类型(需保留原结构时)

如果必须保留Parquet中的嵌套Struct类型,可通过以下方式处理:

  • 在ADF数据集的连接设置中,启用允许读取嵌套类型(部分存储连接器支持该配置)
  • 使用ADF数据流动(Data Flow)替代复制活动,数据流动对嵌套Struct类型的兼容性更好,可通过派生列、扁平化等操作解析嵌套字段

验证步骤

  1. 重新运行修正后的Databricks代码,写入新的Parquet文件
  2. 在ADF中刷新并重新加载数据集,检查字段类型是否均为原始类型(字符串、数值等)
  3. 执行复制活动,确认报错消失

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

相关产品推荐
方舟 Agent Plan

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

最近更新时间:2026.07.11 02:26:09