【问题标题】:spark : use the global config variables in executorsspark : 在执行器中使用全局配置变量
【发布时间】:2015-05-12 04:04:01
【问题描述】:

我的 spark 应用中有一个全局配置对象。

Object Config {
 var lambda = 0.01
}

我会根据用户的输入设置 lambda 的值。

Object MyApp {
   def main(args: String[]) {
     Config.lambda = args(0).toDouble
     ...
     rdd.map(_ * Config.lambda)
   }
}

我发现修改在executors中没有生效。 lambda 的值始终为 0.01。我猜driver的jvm中的修改不会影响executor的。

您还有其他解决方案吗?

我在stackoverflow中发现了一个类似的问题:

how to set and get static variables from spark?

在@DanielL。的回答,他给出了三个解决方案:

  1. 将值放入闭包中,以序列化到执行器以执行任务。

但是我想知道如何编写闭包以及如何将其序列化给执行者,谁能给我一些代码示例?

2.如果值是固定的或配置在执行器节点上可用(位于 jar 内等),那么你可以有一个惰性 val,保证只初始化一次。

如果我将 lambda 声明为惰性 val 变量会怎样?驱动程序中的修改会在执行程序中生效吗?你能给我一些代码示例吗?

3.用数据创建一个广播变量。我知道这种方式,但它还需要一个包装配置对象的本地广播 [] 变量,对吗?例如:

val config = sc.broadcast(Config)

并在 executors 中使用config.value.lambda,对吗?

【问题讨论】:

    标签: apache-spark


    【解决方案1】:
    1. 将值放入闭包中
    object Config {var lambda = 0.01}
    object SOTest {
      def main(args: Array[String]) {
        val sc = new SparkContext(new SparkConf().setAppName("StaticVar"))
        val r = sc.parallelize(1 to 10, 3)
        Config.lambda = 0.02
        mul(r).collect.foreach(println)
        sc.stop()
      }
      def mul(rdd: RDD[Int]) = {
        val l = Config.lambda
        rdd.map(_ * l)
      }
    }
    
    1. lazy val 仅用于初始化一次
    object SOTest {
      def main(args: Array[String]) {
        lazy val lambda = args(0).toDouble
        val sc = new SparkContext(new SparkConf().setAppName("StaticVar"))
        val r = sc.parallelize(1 to 10, 3)
        r.map(_ * lambda).collect.foreach(println)
        sc.stop()
      }
    }
    
    1. 使用数据创建广播变量
    object Config {var lambda = 0.01}
    object SOTest {
      def main(args: Array[String]) {
        val sc = new SparkContext(new SparkConf().setAppName("StaticVar"))
        val r = sc.parallelize(1 to 10, 3)
    
        Config.lambda = 0.04
        val bc = sc.broadcast(Config.lambda)
        r.map(_ * bc.value).collect.foreach(println)
    
        sc.stop()
      }
    }
    

    注意:您不应该将Config Object 直接传递给sc.broadcast(),它会在将您的配置传输给执行程序之前对其进行序列化,但是您的配置是不可序列化的。这里要提到的另一件事:Broadcast variable 不适合您在这里的情况,因为您只共享一个值。

    【讨论】:

    • 非常感谢!!!如果我的配置类有太多变量怎么办?在第一个解决方案中,我应该将 Config 声明为 Class 并新建一个 Config 实例吗?如何将其传递给执行者?
    • 这是正确的吗? def mul(rdd: RDD[Int]) = { val l = new Config() l.lambda=0.02 rdd.map(_ * l.lambda) }
    • @user2848932,如果你想这样做,你应该先让你的配置扩展Serializable
    • 在扩展Serializable之后,我的代码就可以工作了,对吧?
    • 感谢@YijieShen。你能详细说明[2]吗?延迟验证是为 Spark 应用程序初始化一次,还是每个执行程序一次或每个分区一次?您能否对此多说一些?
    猜你喜欢
    • 2017-11-28
    • 2020-01-19
    • 1970-01-01
    • 1970-01-01
    • 2021-08-18
    • 2015-09-04
    • 1970-01-01
    • 1970-01-01
    • 1970-01-01
    相关资源
    最近更新 更多