【问题标题】:How Spark RDD partitions are processed if no. of executors < no. of RDD partition如果没有,如何处理 Spark RDD 分区。执行人<没有。 RDD 分区
【发布时间】:2017-05-03 17:08:02
【问题描述】:

我想了解火花流的基本知识。我有 50 个 Kafka 主题分区和 5 个执行程序,我使用的是 DirectAPI,所以没有。 RDD 分区数将是 50。这个分区将如何在 5 个执行程序上处理?将在每个执行程序上一次触发处理 1 个分区,或者如果执行程序有足够的内存和内核,它将在每个执行程序上并行处理 1 个以上的分区。

【问题讨论】:

    标签: hadoop apache-spark apache-kafka spark-streaming


    【解决方案1】:

    将在每个执行程序上一次触发进程 1 个分区,或者如果 executor 有足够的内存和内核,它将处理超过 1 个 在每个执行器上并行分区。

    Spark 将根据您正在运行的作业可用的内核总数来处理每个分区。

    假设您的流式传输作业有 10 个执行程序,每个执行程序有 2 个内核。这意味着您将能够同时处理 10 x 2 = 20 个分区,假设 spark.task.cpus 设置为 1。

    如果你真的想要详细信息,请查看 Spark Standalone 向CoarseGrainedSchedulerBackend 请求资源,你可以查看它是makeOffers

    private def makeOffers() {
      // Filter out executors under killing
      val activeExecutors = executorDataMap.filterKeys(executorIsAlive)
      val workOffers = activeExecutors.map { case (id, executorData) =>
        new WorkerOffer(id, executorData.executorHost, executorData.freeCores)
      }.toIndexedSeq
      launchTasks(scheduler.resourceOffers(workOffers))
    }
    

    这里的关键是executorDataMap,它保存了从执行程序 id 到 ExecutorData 的映射,它告诉系统中每个这样的执行程序正在使用多少内核,并根据它和分区的首选位置,使有根据地猜测该任务应该运行哪个执行器。

    这是一个使用 Kafka 的实时 Spark Streaming 应用程序的示例:

    我们有 5 个分区,运行 3 个执行器,每个执行器有超过 2 个核心,这使得流能够同时处理每个分区。

    【讨论】:

    • 非常感谢您提供如此精确而详尽的答案,这意味着如果 spark.task.cpus 设置为 1,则每个分区由 1 个核心(1 个线程)处理。所以在我的情况下,如果我有 5 executors 和我设置--executor-cores 10 将同时处理所有分区。
    • @nilesh1212 任务数量取决于DirectKafkaInputDStream 中的分区数,但基本上每个任务都会为底层DStream 的RDD 中的每个分区请求(参见this question for more) .您可以通过进入从 Kafka 读取数据并查看每个分区正在处理的位置的第一个转换来自己验证这一点。
    • 如果我没记错的话,任务只不过是在分区级别的 rdd 上执行的转换/操作?
    • 我认为“运行使用 CoarseGrainedSchedulerBackend 的 Spark Standalone”并不成立。情况正好相反,即 CGSB 用于与 Spark Standalone 和 Mesos 通信。
    • @Jacek 也许“使用”这个词不合适。也许“沟通”更好。
    猜你喜欢
    • 1970-01-01
    • 1970-01-01
    • 2021-02-24
    • 1970-01-01
    • 2019-06-16
    • 1970-01-01
    • 1970-01-01
    • 2017-02-17
    • 2020-09-18
    相关资源
    最近更新 更多