Databricks Notebook中Schema自动推断失败问题排查
我在Databricks中编写了Spark结构化流代码,逻辑为:先检查实体对应的Delta表是否存在,若不存在则创建该表,且希望通过inferSchema选项自动推断Delta表的Schema。
代码实现
# Check if the Delta table exists, if not create it if not DeltaTable.isDeltaTable(spark, sink_path): # Read the Parquet data and infer the schema parquet_data = spark.read.option("inferSchema", "true").parquet(source_path) # Create a Delta table with the inferred schema #### does not create transaction log parquet_data.write.format("delta").mode("overwrite").save(sink_path) print('delta table created')
遇到的错误
batch stream failed for advertiser: An error occurred while calling o671.load.
: com.databricks.sql.cloudfiles.errors.CloudFilesException: Cannot infer schema when the input pathdbfs:/mnt/raw/Entity1/fileA.parquetis empty. Please try to start the stream when there are files in the input path, or specify the schema.
源文件包含约50条记录,请问为何Schema推断无法正常工作?
问题分析与解决办法
可能的原因
源路径或文件异常
- 虽确认源文件有数据,但Spark可能无法正常访问:比如路径拼写错误、云存储权限不足、Parquet文件本身损坏(格式无效)。
- 部分云存储路径区分大小写,需检查
source_path的大小写与实际路径完全匹配。
代码逻辑漏洞
- 当前代码中
parquet_data.write语句不在if分支内:若Delta表已存在,parquet_data变量未定义会触发报错;若进入if分支仍报错,说明读取Parquet时确实未识别到有效数据。
- 当前代码中
CloudFiles组件限制
- 错误提示来自
CloudFilesException,说明你可能在结构化流中使用了CloudFiles组件,它对Schema推断的要求更严格,即使文件非空,格式不符合预期也会被判定为"空路径"。
- 错误提示来自
解决步骤
验证源文件可用性
在Databricks Notebook中执行以下命令,确认能正常读取数据:df = spark.read.parquet("dbfs:/mnt/raw/Entity1/fileA.parquet") print(f"数据条数: {df.count()}") df.printSchema()若执行失败,优先排查路径、权限或文件损坏问题。
修复代码逻辑
将创建Delta表的代码移入if分支,避免未定义变量的问题:
from delta.tables import DeltaTable
if not DeltaTable.isDeltaTable(spark, sink_path):
parquet_data = spark.read.option("inferSchema", "true").parquet(source_path)
parquet_data.write.format("delta").mode("overwrite").save(sink_path)
print('delta table created')
else:
print('delta table already exists')
3. **显式指定Schema(备选方案)** 若自动推断仍失败,显式定义Schema可彻底解决问题: ```python from pyspark.sql.types import StructType, StructField, StringType, IntegerType # 方式1:从现有文件读取Schema sample_df = spark.read.parquet("dbfs:/mnt/raw/Entity1/fileA.parquet") custom_schema = sample_df.schema # 方式2:手动编写Schema # custom_schema = StructType([ # StructField("col1", StringType(), True), # StructField("col2", IntegerType(), True) # ]) if not DeltaTable.isDeltaTable(spark, sink_path): parquet_data = spark.read.schema(custom_schema).parquet(source_path) parquet_data.write.format("delta").mode("overwrite").save(sink_path) print('delta table created')
内容的提问来源于stack exchange,提问作者Shoaib Maroof

