【问题标题】:Can Spark/EMR read data from s3 multi-threadedSpark/EMR 可以从 s3 多线程读取数据吗
【发布时间】:2020-01-20 15:47:11
【问题描述】:

由于一系列不幸的事件,我们最终在 s3 上存储了一个非常碎片化的数据集。表元数据存储在 Glue 上,数据使用“bucketBy”写入,并以 parquet 格式存储。因此,文件的发现不是问题,spark 分区的数量等于桶的数量,这提供了良好的并行度。

当我们在 Spark/EMR 上加载这个数据集时,我们最终会让每个 spark 分区从 s3 加载大约 8k 个文件。

由于我们以列格式存储数据;根据我们需要几个字段的用例,我们并没有真正读取所有数据,而是读取存储的一小部分数据。

根据工作节点上的 CPU 利用率,我可以看到每个任务(每个分区运行)几乎使用了大约 20% 的 CPU,我怀疑这是由于每个任务有一个线程顺序从 s3 读取文件,这么多 IOwait...

有没有办法鼓励 EMR 上的 spark 任务多线程从 s3 读取数据,这样我们就可以在一个任务中同时从 s3 读取多个文件?这样,我们可以利用 80% 的空闲 CPU 来加快速度吗?

【问题讨论】:

  • 一般情况下,一个vcpu每个任务只能运行1个线程。如果你想为每个任务运行多个线程,你需要设置这个变量:spark.task.cpu,默认为 1,然后你需要在你的代码中进行并行处理,这将在 executor 上运行。
  • 这里有一个链接解释了为什么在做 IO 时 cpu 没有被充分利用,stackoverflow.com/questions/13596997/…
  • 我认为我可以在 RDD 级别上执行这样的操作,但是我正在处理 Dataset/Dataframe API 级别,因此我需要一个底层级别来读取多线程文件......除非我最终使用内置的多线程从头开始重新实现 Dataframes ;)
  • 那么你需要创建多个可以并行读取数据的执行器。这完全取决于您如何编写代码。
  • 我不太明白你的意思。能举个例子吗?

标签: multithreading apache-spark amazon-s3 amazon-emr


【解决方案1】:

使用 Spark 数据帧读取 S3 数据有两个部分:

  1. 发现(列出 S3 上的对象)
  2. 读取S3对象,包括解压等

发现通常发生在驱动程序上。一些托管 Spark 环境具有使用集群资源进行更快发现的优化。除非您获得超过 100K 个对象,否则这通常不是问题。如果您有 .option("mergeSchema", true),则发现速度会较慢,因为必须触摸每个文件才能发现其架构。

读取 S3 文件是执行操作的一部分。读取的并行度为 min(分区数,可用内核数)。更多分区 + 更多可用内核意味着更快的 I/O……理论上。实际上,如果您没有定期访问这些文件以使 S3 扩大其可用性,那么 S3 可能会非常慢。因此,在实践中,额外的 Spark 并行性收益递减。观察每个活动核心的总网络 RW 带宽,并调整您的执行以获得最高值。

您可以通过df.rdd.partitions.length查看分区数。

如果 S3 I/O 吞吐量较低,您可以执行其他操作:

  1. 确保 S3 上的数据在其前缀方面是分散的(请参阅https://docs.aws.amazon.com/AmazonS3/latest/dev/optimizing-performance.html)。

  2. 打开 AWS 支持请求并请求扩展数据的前缀。

  3. 用不同的节点类型进行实验。我们发现存储优化节点具有更好的有效 I/O。

希望这会有所帮助。

【讨论】:

  • 嗨@Sim,我要问的是,每个任务如何在多线程中读取(如果可能的话)。我要解决的问题是,每个任务(即每个分区)只使用大约 20% 的 CPU,而 IOWait 是由于从 s3 读取数据。
  • @zetaprime 目前无法使用 Spark 本地执行此操作。它从根本上违反了处理模型。一些托管 Spark 环境具有缓存优化,可以从 S3 到本地工作程序或更快的网络存储进行预取操作,从而实现与您期望的效果相似。独立地,您如何确定瓶颈不在 S3 或本地物理硬件/网络上?对此进行测试的唯一方法是更改​​ Spark 并行性,看看会发生什么。
  • 我可以清楚地看到当我检查工作线程时,在从 s3 请求数据时看到等待套接字读取的线程。
  • @zetaprime 我明白了。我的答案是:除非您编写自定义代码来处理数据并且只使用 Spark 进行并行代码执行,否则如果不增加工作内核或使用带有智能缓存层的托管 Spark 提供程序,就无法提高并行性。
  • 每个进程有多少个工人?确保线程和 http 连接的 s3a 设置大于此值,这样这些池就不会成为瓶颈。否则,小文件到处都是昂贵的,在 S3 上由于开销而更糟。你能在做其他事情之前合并它们吗?
【解决方案2】:

I/O 似乎是特定任务的界限。因此,如果 I/O 吞吐量相对较低,我认为这是由于网络 bindwidth 而不是 cpu 使用率低。

一种可能满足您需求的方法是使用 spark-sql 数据源 api。并在分区阅读器中自定义您的多线程阅读策略。

【讨论】:

  • 您能否解释一下定制或提供此类策略示例的链接?
猜你喜欢
  • 1970-01-01
  • 1970-01-01
  • 2018-08-08
  • 2020-05-19
  • 2021-09-13
  • 1970-01-01
  • 2018-11-09
  • 1970-01-01
  • 2018-09-04
相关资源
最近更新 更多