【问题标题】:How to have idempotent guarantee when writing spark dataset.to hdfs?将spark dataset.写入hdfs时如何获得幂等保证?
【发布时间】:2021-01-12 15:05:33
【问题描述】:

我有一个写入 hdfs(镶木地板文件)的 spark 进程。我的猜测是,默认情况下,如果 spark 有一些失败并重试,它可能会写入两次文件(我错了吗?)。

但是,我该怎么做才能在 hdfs 输出上获得幂等性呢?

我看到 2 种情况应该以不同的方式提出质疑(但请纠正我或如果您了解得更好,请进一步发展):

  1. 写入一项时发生故障:我猜写入已重新启动,因此如果 hdfs 上的帖子不是“原子”w.r.t 以引发写入调用,则可能会写入两次。机会有多大?
  2. 失败发生在任何地方,但是由于执行 dag 的制作方式,重新启动将发生在几个写入任务之前的任务中(例如,我正在考虑必须在某些 groupBy 之前重新启动),还有一些这些写入任务中的一部分已经完成。 Spark 执行是否保证不会再次调用这些任务?

【问题讨论】:

    标签: apache-spark hdfs idempotent


    【解决方案1】:

    我认为这取决于您在工作中使用什么样的提交者,以及该提交者是否能够撤消失败的工作。例如 当您使用 Apache Parquet 格式的输出时,Spark 期望提交者 Parquet 是 ParquetOutputCommitter 的子类。如果您使用此提交者DirectParquetOutputCommitter 来附加数据,则无法撤消该作业。 code

    如果您使用ParquetOutputCommitter 本身,您可以see 扩展FileOutputCommitter 并稍微覆盖commitJob(JobContext jobContext) 方法。

    以下内容从Hadoop: The Definitive Guide复制/粘贴

    OutputCommitter API: setupJob() 方法在作业运行之前调用,通常用于执行 初始化。对于FileOutputCommitter,该方法创建最终输出目录, ${mapreduce.output.fileoutputformat.outputdir},还有一个临时工作空间 对于任务输出,_temporary,作为它下面的子目录。 如果作业成功,则调用 commitJob() 方法,该方法在默认的基于文件的 实现删除临时工作空间并创建一个隐藏的空标记 输出目录中名为 _SUCCESS 的文件,以向文件系统客户端指示该作业 顺利完成。如果作业未成功,则使用状态对象调用 abortJob() 指示作业是失败还是被杀死(例如,被用户杀死)。在默认 实施,这将删除作业的临时工作空间。

    任务级别的操作类似。 setupTask() 方法在 任务运行,默认实现不做任何事情,因为临时 为任务输出命名的目录是在写入任务输出时创建的。

    任务的提交阶段是可选的,可以通过返回 false 来禁用 需要任务提交()。这使框架不必运行分布式提交 该任务的协议,并且既不调用 commitTask() 也不调用 abortTask()。 FileOutputCommitter 将在未写入任何输出时跳过提交阶段 任务。

    如果任务成功,commitTask() 会被调用,在默认实现中会移动 临时任务输出目录(其名称中包含任务尝试 ID,以避免 任务尝试之间的冲突)到最终输出路径, ${mapreduce.output.fileoutputformat.outputdir}。否则,框架调用 abortTask(),删除临时任务输出目录。

    该框架确保在针对特定任务进行多次任务尝试的情况下, 只有一个会被承诺;其他的将被中止。出现这种情况可能是因为 第一次尝试由于某种原因失败了——在这种情况下,它会被中止,然后, 成功的尝试将被提交。如果两次任务尝试 作为投机副本同时运行;在这种情况下,第一个完成的 将被提交,而另一个将被中止。

    【讨论】:

      猜你喜欢
      • 1970-01-01
      • 2020-09-15
      • 2018-06-02
      • 1970-01-01
      • 1970-01-01
      • 2013-06-18
      • 2018-12-25
      • 2015-12-10
      • 1970-01-01
      相关资源
      最近更新 更多