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)
当您每天进行压缩时,分区变量应该代表昨天的数据分区。没有新文件写入该分区,并且增量日志可以识别并发更改将不适用于今天当前传入的数据。
| 归档时间: |
|
| 查看次数: |
13304 次 |
| 最近记录: |