更新:
如果您使用 AWS Elastic MapReduce,版本 >= 5.19 的集群现在可以安全地使用推测执行,但您的 Spark 作业仍然可能中途失败并留下不完整的结果。
如果您直接扫描 AWS S3,您的不完整结果中的部分数据是可查询的,这可能会导致下游作业的结果不正确,因此您需要一种策略来处理它!
如果您运行的是 Spark 2.3.0 或更高版本,我建议您使用 SaveMode.Overwrite 将新分区写入确定性位置并在失败时重试,这样可以避免输出中出现重复或损坏的数据。
如果您使用的是SaveMode.Append,那么重试 Spark 作业会在输出中产生重复数据。
推荐的方法:
df.write
.mode(SaveMode.Overwrite)
.partitionBy("date")
.parquet("s3://myBucket/path/to/table.parquet")
然后在成功写入分区后,以原子方式将其注册到 Hive 等元存储,并查询 Hive 作为您的真实来源,而不是直接查询 S3。
例如。
ALTER TABLE my_table ADD PARTITION (date='2019-01-01') location 's3://myBucket/path/to/table.parquet/date=2019-01-01'
如果您的 Spark 作业失败并且您正在使用 SaveMode.Overwrite,那么重试始终是安全的,因为数据尚未可用于 Metastore 查询,而您只是覆盖了失败分区中的数据。
注意:为了只覆盖特定分区而不是您需要配置的整个数据集:
spark.conf.set("spark.sql.sources.partitionOverwriteMode","dynamic")
仅适用于 Spark 2.3.0。
https://aws.amazon.com/blogs/big-data/improve-apache-spark-write-performance-on-apache-parquet-formats-with-the-emrfs-s3-optimized-committer/
https://docs.aws.amazon.com/emr/latest/ReleaseGuide/emr-spark-s3-optimized-committer.html
随着 Iceberg 项目的成熟,您可能还想将其视为 Hive / Glue 元存储的替代方案。
https://github.com/apache/incubator-iceberg
为什么这是必要的以及非 AWS 用户的背景
在提交到对象存储时使用 Spark 推测通常是一个非常糟糕的主意,具体取决于在下游查看该数据的内容和您的一致性模型。
来自 Netflix 的 Ryan Blue 进行了一场精彩(而且非常有趣)的演讲,准确地解释了原因:https://www.youtube.com/watch?v=BgHrff5yAQo
从 OP 的描述来看,我怀疑他们正在写 Parquet。
TL;dr 版本是在 S3 中,重命名操作实际上是复制和删除,这具有一致性含义。通常在 Spark 中,输出数据被写入临时文件位置,并在计算完成后重命名。这意味着如果投机执行处于开启状态,那么多个执行器可以处理相同的结果,然后通过将临时文件重命名为最终结果,第一个完成的执行器“获胜”,而另一个任务被取消。此重命名操作发生在单个任务上,以确保只有一个推测性任务获胜,这在 HDFS 上不是问题,因为重命名是一种廉价的元数据操作,其中几千或几百万次只需要很少的时间。
但是在使用 S3 时,重命名不是原子操作,它实际上是一个需要时间的副本。因此,您可能会遇到一种情况,即您必须第二次在 S3 中复制一大堆文件以连续重命名,这是一个同步操作,导致您的速度变慢。如果您的执行程序有多个核心,您实际上可能有一个任务破坏了另一个任务的结果,这在理论上应该没问题,因为一个文件最终获胜,但您无法控制此时发生的事情。
问题是,如果最终的重命名任务失败会发生什么?您最终会将部分文件提交给 S3,而不是全部提交,这意味着部分/重复的数据以及下游的许多问题,具体取决于您的应用程序。
虽然我不喜欢,但目前流行的做法是在本地写入 HDFS,然后使用 S3Distcp 之类的工具上传数据。
看看 HADOOP-13786。
Steve Loughran 是这个问题的最佳人选。
如果您不想等待,Ryan Blue 有一个 repo “rdblue/s3committer”,它允许您为除 parquet 文件之外的所有输出修复此问题,但正确集成和子类化似乎需要一些工作。
更新:
HADOOP-13786 现已修复并发布到 Hadoop 3.1 库中。
目前 Steven Loughran 正在努力获得一个基于 Hadoop 3.1 库并入 apache/spark (SPARK-23977) 的修复,但是根据票务评论线程的最新消息是,该修复不会在 Spark 2.4 发布之前被合并,所以我们可以等待它成为主流。
更新 v2:
注意:您可以通过在 Hadoop 配置中将 mapreduce.fileoutputcommitter.algorithm.version 设置为 2 来将最终输出分区重命名任务可能失败的时间窗口减半,因为原始输出提交机制实际上执行了 两次 重命名.