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

Databricks中CSV导入Delta表遇DELTA_INVALID_FORMAT错误求助

问题:Databricks中CSV导入Delta表时出现格式不兼容错误

将CSV文件导入Delta表的初始加载流程拆分为3个Notebook,前两步已执行成功,第三步加载CSV并归档时触发格式不兼容错误。


已完成的操作步骤

1. 挂载Blob存储

#mount the BLOB storage in databricks
dbutils.fs.mount(
    source="wasbs://CONTAINORXX@databrickstoragetest.blob.core.windows.net/",
    mount_point="/mnt/MOUNTXX",
    extra_configs={"fs.azure.account.key.databrickstoragetest.blob.core.windows.net": "KEYXX"}
)

2. 创建Schema和Delta表

#create delta table
from delta.tables import DeltaTable
from pyspark.sql.types import StructType, StructField, StringType, IntegerType, DateType

# create schema
schema = StructType([
    StructField("name", StringType(), True),
    StructField("age", IntegerType(), True),
    StructField("country", StringType(), True)
])

#create the new table
empty_df = spark.createDataFrame([], schema)
empty_df.write.format("delta").saveAsTable("persontest")

出错的第三步代码及错误信息

执行代码

# Define file paths
input_path = "/mnt/MOUNTXX/*.csv"
archive_path = "/mnt/MOUNTXX/archive"

#read csv
df = spark.read.format("csv").option("header", "true").option("inferSchema", "true").load(input_path)

#load data into table
df.write.format("delta").mode("append").saveAsTable("persontest")

#archive file
processed_files = dbutils.fs.ls(input_path)

for file in processed_files:
    if file.path.endswith(".csv"):
        dbutils.fs.mv(file.path, archive_path + "/" + file.name)

错误信息

[DELTA_INVALID_FORMAT] Incompatible format detected. 
A transaction log for Delta was found at `/_delta_log`, 
but you are trying to read from `/mnt/MOUNTXX/*.csv` using format("csv"). 
You must use 'format("delta")' when reading and writing to a delta table. 
SQLSTATE: 22000
File <command-2260341893915879>, line 6      
3 archive_path = "/mnt/MOUNTXX/archive"      
5 #read csv ----> 
6 df = spark.read.format("csv").option("header", "true").option("inferSchema", "true").load(input_path)      
8 #load data into table      
9 df.write.format("delta").mode("append").saveAsTable("persontest") 
File /databricks/spark/python/pyspark/instrumentation_utils.py:47, in _wrap_function.<locals>.wrapper(*args, **kwargs)     
45 start = time.perf_counter()     
46 try: ---> 
47     res = func(*args, **kwargs)    
 48     logger.log_success(     
49         module_name, class_name, function_name, time.perf_counter() - start, signature     
50     )     
51     return res 
File /databricks/spark/python/pyspark/sql/readwriter.py:312, in DataFrameReader.load(self, path, format, schema, **options)    
310 self.options(**options)    
311 if isinstance(path, str): ---> 
312     return self._df(self._jreader.load(path))    
313 elif path is not None:    
314     if type(path) != list: 
File /databricks/spark/python/lib/py4j-0.10.9.7-src.zip/py4j/java_gateway.py:1355, in JavaMember.__call__(self, *args)   
1349 command = proto.CALL_COMMAND_NAME +\   
1350     self.command_header +\   
1351     args_command +\   
1352     proto.END_COMMAND_PART   
1354 answer = self.gateway_client.send_command(command) ---> 
1355 return_value = get_return_value(   
1356     answer, self.gateway_client, self.target_id, self.name)   
1358 for temp_arg in temp_args:   
1359     if hasattr(temp_arg, "_detach"): 
File /databricks/spark/python/pyspark/errors/exceptions/captured.py:261, in capture_sql_exception.<locals>.deco(*a, **kw)    
257 converted = convert_exception(e.java_exception)    
258 if not isinstance(converted, UnknownException):    
259     # Hide where the exception came from that shows a non-Pythonic    
260     # JVM exception message. ---> 
261     raise converted from None    
262 else:    
263     raise

问题原因

Spark在/mnt/MOUNTXX路径下检测到Delta事务日志目录/_delta_log,因此将该路径识别为Delta表存储路径,拒绝使用CSV格式读取该路径下的文件。这通常是因为:

  • 之前误将Delta表写入到挂载根目录
  • 根目录下残留了其他Delta操作生成的日志文件

解决方案

方案1:使用独立的CSV输入目录(推荐)

在Blob容器内创建专门的input子目录存放CSV文件,确保该目录下无Delta日志文件,修改代码如下:

# 导入之前定义的schema(如果在当前Notebook未定义,需重新导入)
from pyspark.sql.types import StructType, StructField, StringType, IntegerType, DateType

schema = StructType([
    StructField("name", StringType(), True),
    StructField("age", IntegerType(), True),
    StructField("country", StringType(), True)
])

# Define file paths
input_path = "/mnt/MOUNTXX/input/*.csv"
archive_path = "/mnt/MOUNTXX/archive"

# 读取CSV时使用预定义schema,避免inferSchema导致的类型不匹配
df = spark.read.format("csv").option("header", "true").schema(schema).load(input_path)

# 写入Delta表
df.write.format("delta").mode("append").saveAsTable("persontest")

# 归档文件
processed_files = dbutils.fs.ls("/mnt/MOUNTXX/input")

for file in processed_files:
    if file.path.endswith(".csv"):
        dbutils.fs.mv(file.path, f"{archive_path}/{file.name}")

优势:从根源上隔离CSV文件与Delta表存储,避免后续操作再次触发格式冲突。

方案2:清理根目录下的Delta日志(谨慎操作)

如果确认/mnt/MOUNTXX/_delta_log是误生成的冗余日志,可执行以下命令删除:

dbutils.fs.rm("/mnt/MOUNTXX/_delta_log", recurse=True)

注意:删除前务必确认该日志不属于任何正在使用的Delta表,避免数据丢失或损坏。


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

相关产品推荐
方舟 Agent Plan

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

最近更新时间:2026.06.15 18:24:50