【问题标题】:getting number of visible nodes in PySpark获取 PySpark 中可见节点的数量
【发布时间】:2015-04-30 09:01:22
【问题描述】:

我正在 PySpark 中运行一些操作,并且最近增加了我的配置(在 Amazon EMR 上)中的节点数量。然而,即使我将节点数量增加了两倍(从 4 个到 12 个),性能似乎并没有改变。因此,我想看看新节点是否对 Spark 可见。

我正在调用以下函数:

sc.defaultParallelism
>>>> 2

但我认为这告诉我分配给每个节点的任务总数,而不是 Spark 可以看到的节点总数。

如何查看 PySpark 在我的集群中使用的节点数量?

【问题讨论】:

    标签: python-2.7 apache-spark pyspark


    【解决方案1】:

    在 pyspark 上,您仍然可以使用 pyspark 的 py4j 桥接器调用 scala getExecutorMemoryStatus API:

    sc._jsc.sc().getExecutorMemoryStatus().size()
    

    【讨论】:

    • 出于某种原因,这似乎对我不起作用。我发布了一个带有最小示例的question,包括输出(我从这个调用中得到 1,而实际上有 12 个执行者/工人)。
    【解决方案2】:

    sc.defaultParallelism 只是一个提示。根据配置,它可能与节点数量无关。如果您使用带有分区计数参数但您不提供它的操作,则这是分区数。例如sc.parallelize 将从列表中创建一个新的 RDD。您可以使用第二个参数告诉它要在 RDD 中创建多少个分区。但是这个参数的默认值是sc.defaultParallelism

    您可以在 Scala API 中使用 sc.getExecutorMemoryStatus 获取执行器的数量,但这在 Python API 中没有公开。

    一般而言,建议是 RDD 中的分区数量大约是执行程序数量的 4 倍。这是一个很好的提示,因为如果任务花费的时间有差异,这将平衡它。例如,一些 executor 将处理 5 个更快的任务,而另一些 executor 将处理 3 个较慢的任务。

    您不需要非常准确。如果您有一个粗略的想法,您可以进行估算。就像如果你知道你有少于 200 个 CPU,你可以说 500 个分区就可以了。

    所以尝试用这个分区数创建 RDD:

    rdd = sc.parallelize(data, 500)     # If distributing local data.
    rdd = sc.textFile('file.csv', 500)  # If loading data from a file.
    

    如果您不控制 RDD 的创建,或者在计算之前重新分区 RDD:

    rdd = rdd.repartition(500)
    

    您可以使用rdd.getNumPartitions()查看RDD中的分区数。

    【讨论】:

    • 谢谢。但是,当我运行sc.getExecutorMemoryStatus 时,我收到一条错误消息'SparkContext' object has no attribute 'getExecutorMemoryStatus'。你在使用 PySpark 吗?
    • 不,我使用的是 Scala API。我以为它也是 Python API 的一部分,但似乎不是。我想您的选择是自己将其添加到 Python API 或切换到 Scala API。或者将机器数量作为命令行参数。
    • 将机器数量作为命令行参数是什么意思?打开 PySpark 交互式 shell 时可以调用它吗?我是一个 Python 人(所以没有 Scala)并且不想开始进行 API 更改。
    • 在创建 RDD 时设置它的分区数。只要您的分区数量多于执行程序核心的数量,所有执行程序都会有一些工作要做。所以确切的计数并不那么重要。您可以使用rdd.getNumPartitions() 查看 RDD 中的分区数。或者使用rdd.repartition(n)改变分区数(这是一个shuffle操作)。
    • 这个答案并不能真正回答问题,您可以通过 pyspark 访问getExecutorMemoryStatus
    【解决方案3】:

    使用这个应该可以获取集群中的节点数(类似于上面@Dan的方法,但是更短,效果更好!)。

    sc._jsc.sc().getExecutorMemoryStatus().keySet().size()
    

    【讨论】:

      【解决方案4】:

      其他答案提供了一种获取执行者数量的方法。这是一种获取节点数量的方法。这包括头节点和工作节点。

      s = sc._jsc.sc().getExecutorMemoryStatus().keys()
      l = str(s).replace("Set(","").replace(")","").split(", ")
      
      d = set()
      for i in l:
          d.add(i.split(":")[0])
      len(d)  
      

      【讨论】:

        【解决方案5】:

        我发现有时我的会话被远程杀死了一个奇怪的 Java 错误

        Py4JJavaError: An error occurred while calling o349.defaultMinPartitions.
        : java.lang.IllegalStateException: Cannot call methods on a stopped SparkContext.
        

        我通过以下方式避免了这种情况

        def check_alive(spark_conn):
            """Check if connection is alive. ``True`` if alive, ``False`` if not"""
            try:
                get_java_obj = spark_conn._jsc.sc().getExecutorMemoryStatus()
                return True
            except Exception:
                return False
        
        def get_number_of_executors(spark_conn):
            if not check_alive(spark_conn):
                raise Exception('Unexpected Error: Spark Session has been killed')
            try:
                return spark_conn._jsc.sc().getExecutorMemoryStatus().size()
            except:
                raise Exception('Unknown error')
        

        【讨论】:

        • 机器数不等于执行器数。一台机器有一个或多个执行者。
        • 谢谢这是一个错字。我的目的只是为了避免会话被杀死给出错误
        猜你喜欢
        • 1970-01-01
        • 1970-01-01
        • 1970-01-01
        • 2017-09-04
        • 1970-01-01
        • 2018-05-30
        • 1970-01-01
        • 1970-01-01
        • 1970-01-01
        相关资源
        最近更新 更多