启用Unity环境下Databricks AutoLoader读取CSV失败问题
问题:Databricks AutoLoader无法从非空CSV文件推断Schema
环境背景
- 已启用Unity Catalog
- 操作目标:通过Databricks AutoLoader读取CSV文件并保存为Delta表
问题现象
执行AutoLoader流式代码时,系统报错称输入路径为空,但实际CSV文件非空——批处理读取代码可正常生成Delta表,也能通过dbutils命令查看文件内容。
AutoLoader流式处理代码
df_current = spark.readStream.format("cloudFiles") \ .option("cloudFiles.format", "csv") \ .option("multiLine", "true")\ .option("header", "true")\ .option('escape', "\"")\ .option("cloudFiles.schemaLocation", f"{base_loc_out}/{channel_name}/schema/{filename}")\ .option("cloudFiles.inferColumnTypes", "true")\ .load(channel_file.path)\ .writeStream \ .format("delta")\ .option("checkpointLocation", f"{base_loc_out}/{channel_name}/checkpoint/{filename}") \ .trigger(availableNow=True)\ .toTable(f"{catalog_name}.{schema_name}.{filename}") df_current.awaitTermination() print(f"number of rows read => {df_current.lastProgress['numInputRows']}, number of rows written => {df_current.lastProgress['sink']['numOutputRows']}")
错误信息
com.databricks.sql.cloudfiles.errors.CloudFilesException: Cannot infer schema when the input path `abfss://***************/employees.csv` is empty. Please try to start the stream when there are files in the input path, or specify the schema.
可正常运行的批处理读取代码
df_current = spark.read\ .option("header","true")\ .option("inferSchema","true")\ .option("multiLine", "true")\ .option('escape', "\"")\ .option('delimiter', ',')\ .csv(channel_file.path) df_current.write.format("delta").mode("overwrite").saveAsTable(f"{catalog_name}.{schema_name}.{filename}")
文件验证命令
dbutils.fs.head("abfss://****************/employees.csv")
CSV文件内容
empid,empName,salary,dep 1,mani,6000,IT 2,Bharath,2000,Payroll 3,Vamsi,3000,IT 4,Nikhil,1000,Payroll 5,mahesh,2500,HR 6,srinu,3500,Payroll 7,nag,4500,IT 8,neeraj,1500,HR 9,dsp,10000,HR 10,deeraj,4500,IT
解决方案建议
1. 显式指定Schema(最直接解决方法)
绕过AutoLoader的自动推断逻辑,手动定义匹配CSV结构的Schema:
from pyspark.sql.types import StructType, StructField, IntegerType, StringType schema = StructType([ StructField("empid", IntegerType(), nullable=True), StructField("empName", StringType(), nullable=True), StructField("salary", IntegerType(), nullable=True), StructField("dep", StringType(), nullable=True) ]) df_current = spark.readStream.format("cloudFiles") \ .option("cloudFiles.format", "csv") \ .option("multiLine", "true")\ .option("header", "true")\ .option('escape', "\"")\ .option("cloudFiles.schemaLocation", f"{base_loc_out}/{channel_name}/schema/{filename}")\ .schema(schema)\ .load(channel_file.path)\ .writeStream \ .format("delta")\ .option("checkpointLocation", f"{base_loc_out}/{channel_name}/checkpoint/{filename}") \ .trigger(availableNow=True)\ .toTable(f"{catalog_name}.{schema_name}.{filename}") df_current.awaitTermination()
2. 检查路径与权限
- 确认
channel_file.path精确指向单个CSV文件,而非目录(AutoLoader加载目录时会扫描所有文件,但单个文件路径需确保无拼写错误) - 验证Unity Catalog下的存储凭据权限,确保AutoLoader拥有读取该ABFS路径的完整权限(包括文件元数据读取)
3. 清理Schema缓存目录
如果cloudFiles.schemaLocation指定的目录中存在之前生成的空推断记录,删除该目录后重新运行任务,避免旧缓存干扰Schema推断。
4. 补充delimiter参数
在AutoLoader配置中显式添加delimiter=','参数,与批处理代码保持一致,避免默认分隔符匹配问题:
df_current = spark.readStream.format("cloudFiles") \ .option("cloudFiles.format", "csv") \ .option("multiLine", "true")\ .option("header", "true")\ .option('escape', "\"")\ .option('delimiter', ',')\ # 新增参数 .option("cloudFiles.schemaLocation", f"{base_loc_out}/{channel_name}/schema/{filename}")\ .option("cloudFiles.inferColumnTypes", "true")\ .load(channel_file.path)\ ... # 后续写入逻辑不变
内容的提问来源于stack exchange,提问作者Rakesh Prasad
相关产品推荐
相关产品推荐

