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类型的兼容性更好,可通过派生列、扁平化等操作解析嵌套字段
验证步骤
- 重新运行修正后的Databricks代码,写入新的Parquet文件
- 在ADF中刷新并重新加载数据集,检查字段类型是否均为原始类型(字符串、数值等)
- 执行复制活动,确认报错消失
内容的提问来源于stack exchange,提问作者Developer Rajinikanth
相关产品推荐
相关产品推荐

