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
相关产品推荐
相关产品推荐

