异常:org.apache.spark.sql.delta.ConcurrentAppendException:文件通过并发更新添加到表的根目录

Vla*_*mir 6 parquet spark-streaming databricks delta-lake

我有一个简单的 Spark 作业,将数据流式传输到 Delta 表。该表非常小并且没有分区。

创建了许多小镶木地板文件。

按照文档(https://docs.delta.io/1.0.0/best-practices.html)中的建议,我添加了每天运行一次的压缩作业。

    val path = "..."
    val numFiles = 16
    
    spark.read
     .format("delta")
     .load(path)
     .repartition(numFiles)
     .write
     .option("dataChange", "false")
     .format("delta")
     .mode("overwrite")
     .save(path)
Run Code Online (Sandbox Code Playgroud)

每次压缩作业运行时,流作业都会出现以下异常:

org.apache.spark.sql.delta.ConcurrentAppendException: Files were added to the root of the table by a concurrent update. Please try the operation again.
Run Code Online (Sandbox Code Playgroud)

我尝试将以下配置参数添加到流作业中:

spark.databricks.delta.retryWriteConflict.enabled = true  # would be false by default
spark.databricks.delta.retryWriteConflict.limit = 3  # optionally limit the maximum amout of retries
Run Code Online (Sandbox Code Playgroud)

这没有帮助。

知道如何解决这个问题吗?

小智 2

当您流式传输数据时,将创建小文件(附加),并且这些文件将在您的增量日志中引用(更新)。当您执行压缩时,您会尝试通过将数据整理成较大的文件(当前为 16 个)来解决小文件的开销。这些大文件与小文件一起创建,但在写入增量日志时会发生更改。也就是说,事务 0-100 生成 100 个小文件,发生压缩,新事务告诉您现在改为引用 16 个大文件。问题是,在压缩发生时,流作业中已经发生了事务 101-110。毕竟,您正在压缩所有数据,并且本质上存在合并冲突。

解决方案是进入最佳实践中的下一步,并且仅使用以下方法压缩选择的分区:

.option("replaceWhere", partition)
Run Code Online (Sandbox Code Playgroud)

当您每天进行压缩时,分区变量应该代表昨天的数据分区。没有新文件写入该分区,并且增量日志可以识别并发更改将不适用于今天当前传入的数据。