使用在 YARN 集群模式下运行的 spark 2.4.4 和 spark FIFO 调度程序。
我正在使用具有可变线程数的线程池执行程序提交多个 spark 数据帧操作(即将数据写入 S3)。如果我有大约 10 个线程,这可以正常工作,但是如果我使用数百个线程,则会出现死锁,没有根据 Spark UI 安排作业。
哪些因素控制可以同时调度多少作业?驱动程序资源(例如内存/内核)?其他一些火花配置设置?
编辑:
这是我的代码的简要概述
ExecutorService pool = Executors.newFixedThreadPool(nThreads);
ExecutorCompletionService<Void> ecs = new ExecutorCompletionService<>(pool);
Dataset<Row> aHugeDf = spark.read.json(hundredsOfPaths);
List<Future<Void>> futures = listOfSeveralHundredThings
.stream()
.map(aThing -> ecs.submit(() -> {
df
.filter(col("some_column").equalTo(aThing))
.write()
.format("org.apache.hudi")
.options(writeOptions)
.save(outputPathFor(aThing));
return null;
}))
.collect(Collectors.toList());
IntStream.range(0, futures.size()).forEach(i -> ecs.poll(30, TimeUnit.MINUTES));
exec.shutdownNow();
Run Code Online (Sandbox Code Playgroud)
在某些时候,随着nThreads增加,spark 似乎不再安排任何作业,如下所示:
ecs.poll(...) 最终超时nThreads没有正在运行的作业 ID 的正在运行的查询我的执行环境是
m5.4xlargerd5.24xlargespark.driver.cores=24spark.driver.memory=32gspark.executor.memory=21gspark.scheduler.mode=FIFO如果可能,将作业的输出写入 AWS Elastic MapReduce hdfs(以利用本地 hdfs 的几乎即时重命名和更好的文件 IO),并添加 dstcp 步骤将文件移动到 S3,从而省去处理这些文件的所有麻烦。试图成为文件系统的对象存储的内部结构。此外,写入本地 HDFS 将允许您启用推测来控制失控任务,而不会陷入与 DirectOutputCommiter 相关的死锁陷阱。
如果必须使用 S3 作为输出目录,请确保设置以下 Spark 配置
spark.hadoop.mapreduce.fileoutputcommitter.algorithm.version 2
spark.speculation false
Run Code Online (Sandbox Code Playgroud)
注意:由于数据丢失的可能性,DirectParquetOutputCommitter 已从 Spark 2.0 中删除。不幸的是,在我们提高 S3a 的一致性之前,我们必须使用变通办法。Hadoop 2.8 的情况正在改善
避免按字典顺序使用键名。人们可以使用散列/随机前缀或反向日期时间来解决这个问题。诀窍是分层命名您的键,将您过滤的最常见的内容放在键的左侧。由于 DNS 问题,存储桶名称中永远不要有下划线。
启用fs.s3a.fast.upload upload单个文件的部分并行到 Amazon S3
有关更多详细信息,请参阅这些文章 -
在 Spark 2.1.0 中写入 s3 时设置spark.speculation
| 归档时间: |
|
| 查看次数: |
1337 次 |
| 最近记录: |