【问题标题】:Delta Lake Compacting Multiple files to single fileDelta Lake 将多个文件压缩为单个文件
【发布时间】:2020-02-09 21:26:03
【问题描述】:

我目前正在探索由 databricks 开源的 delta Lake。我正在使用 delta Lake 格式读取 kafka 数据并以流形式写入。 Delta Lake 在来自 kafka 的流式写入过程中创建了许多文件,我觉得这是 hdfs 文件系统。

我已尝试将多个文件压缩为单个文件。

val spark =  SparkSession.builder
    .master("local")
    .appName("spark session example")
    .getOrCreate()

  val df = spark.read.parquet("deltalakefile/data/")

  df.repartition(1).write.format("delta").mode("overwrite").save("deltalakefile/data/")
  df.show()

  spark.conf.set("spark.databricks.delta.retentionDurationCheck.enabled","false")

  DeltaTable.forPath("deltalakefile/data/").vacuum(1)

但是当我检查输出时,它正在创建新文件而不是删除任何现有文件。

有没有办法做到这一点。还有这里的保留期是什么关系?使用的时候我们应该如何在HDFS中配置呢?当我想构建具有 delta Lake 格式的原始/青铜层并且我想长期保存我的所有数据(本地数年/云上无限时间)时,我的保留配置应该是什么?

【问题讨论】:

    标签: databricks delta-lake


    【解决方案1】:

    根据设计,Delta 不会立即删除文件以防止活跃的消费者受到影响。它还提供版本控制(又称时间旅行),因此您可以在必要时查看历史记录。要删除以前的版本或未提交的文件,您需要运行 vacuum

    import io.delta.tables._
    
    val deltaTable = DeltaTable.forPath(spark, pathToTable)
    
    deltaTable.vacuum() // use default retention period
    

    关于如何管理青铜/白银/黄金模型的保留和压实的问题,您应该将登陆表(又名青铜)视为仅附加日志。这意味着您不需要在事后执行压缩或任何重写。青铜表应该是您从上游数据源(例如 Kafka)摄取的数据的记录,并且应用了最少的处理。

    青铜表通常用作增量流源来填充下游数据集。鉴于从 Delta 读取是从事务日志中完成的,与使用执行慢速文件列表的标准文件读取器相比,小文件不是这样的问题。

    但是,当您将文件写入青铜表时,仍有一些选项可以优化文件:1) 通过首先重新分区以减少文件数量,在写入 Delta 时压缩您的 Kafka 消息,2) 增加您的触发间隔,因此摄取运行的频率较低,并将更多消息写入更大的文件。

    【讨论】:

    • 嗨西尔维奥感谢您的帮助。好的,所以如果我有文件让我们说一天大,那么只有我上面的命令才能工作是对的。如果我使用 vaccum(1) 是否会删除所有已提交的旧文件,请告诉我保留期限。?
    • 如果您运行vacuum(1),您的意思是删除超过 1 小时的所有内容。默认保留期为 7 天。
    • 我可以覆盖它并使其删除 1 小时前的文件。
    • 是的,您可以指定 0,但请阅读文档中的警告 docs.delta.io/latest/delta-utility.html#delta-vacuum
    • 是的,通过它。那么压缩文件的其他方法可能是什么。
    猜你喜欢
    • 1970-01-01
    • 2013-09-13
    • 1970-01-01
    • 1970-01-01
    • 2017-07-02
    • 1970-01-01
    • 2021-01-20
    • 1970-01-01
    • 1970-01-01
    相关资源
    最近更新 更多