【问题标题】:Spark - Random Number GenerationSpark - 随机数生成
【发布时间】:2016-04-06 15:03:26
【问题描述】:

我编写了一个必须考虑随机数来模拟伯努利分布的方法。我正在使用random.nextDouble 生成一个介于 0 和 1 之间的数字,然后根据给定我的概率参数的值做出决定。

我的问题是 Spark 在我的 for 循环映射函数的每次迭代中生成相同的随机数。我正在使用DataFrame API。我的代码遵循这种格式:

val myClass = new MyClass()
val M = 3
val myAppSeed = 91234
val rand = new scala.util.Random(myAppSeed)

for (m <- 1 to M) {
  val newDF = sqlContext.createDataFrame(myDF
    .map{row => RowFactory
      .create(row.getString(0),
        myClass.myMethod(row.getString(2), rand.nextDouble())
    }, myDF.schema)
}

这是课程:

class myClass extends Serializable {
  val q = qProb

  def myMethod(s: String, rand: Double) = {
    if (rand <= q) // do something
    else // do something else
  }
}

每次调用myMethod 时,我都需要一个新的随机数。我还尝试使用java.util.Randomscala.util.Random v10 不扩展Serializable)在我的方法中生成数字,如下所示,但我仍然在每个 for 循环中得到相同的数字

val r = new java.util.Random(s.hashCode.toLong)
val rand = r.nextDouble()

我做了一些研究,这似乎与 Sparks 确定性有关。

【问题讨论】:

    标签: scala random apache-spark spark-dataframe


    【解决方案1】:

    只要使用SQL函数rand

    import org.apache.spark.sql.functions._
    
    //df: org.apache.spark.sql.DataFrame = [key: int]
    
    df.select($"key", rand() as "rand").show
    +---+-------------------+
    |key|               rand|
    +---+-------------------+
    |  1| 0.8635073400704648|
    |  2| 0.6870153659986652|
    |  3|0.18998048357873532|
    +---+-------------------+
    
    
    df.select($"key", rand() as "rand").show
    +---+------------------+
    |key|              rand|
    +---+------------------+
    |  1|0.3422484248879837|
    |  2|0.2301384925817671|
    |  3|0.6959421970071372|
    +---+------------------+
    

    【讨论】:

    • 这并没有完全解决我的问题,但它是一个优雅的解决方案,我将来可能会使用,所以 +1
    【解决方案2】:

    根据this post,最好的解决办法是不要把new scala.util.Random放在map里面,也不要完全放在外面(即在驱动代码里),而是放在一个中间的mapPartitionsWithIndex

    import scala.util.Random
    val myAppSeed = 91234
    val newRDD = myRDD.mapPartitionsWithIndex { (indx, iter) =>
       val rand = new scala.util.Random(indx+myAppSeed)
       iter.map(x => (x, Array.fill(10)(rand.nextDouble)))
    }
    

    【讨论】:

    • 必须维护使用此解决方案的代码,并希望与社区分享此解决方案有其缺点并且可能会严重影响您的统计分析,请注意。当您的 rdd 分区>1 时,您的 rdd 随机数序列将为每个分区重新开始,每个分区具有新的种子和不同的数字,但它可能会改变整个序列的“特征”。我的建议:不要使用这种方法。
    • @d-xa 感谢您的评论。您能推荐一种替代方法吗?
    • 如果有人使用这种方法,我建议将 myRDD 的分区修复为 1
    【解决方案3】:

    重复相同序列的原因是随机生成器是在数据分区之前创建并使用种子初始化的。然后每个分区从相同的随机种子开始。也许不是最有效的方法,但以下应该可行:

    val myClass = new MyClass()
    val M = 3
    
    for (m <- 1 to M) {
      val newDF = sqlContext.createDataFrame(myDF
        .map{ 
           val rand = scala.util.Random
           row => RowFactory
          .create(row.getString(0),
            myClass.myMethod(row.getString(2), rand.nextDouble())
        }, myDF.schema)
    }
    

    【讨论】:

    • 我稍作修改以解决我的问题。我将 Random val 传递到我的方法中,并从那里生成随机数。这解决了我的问题,但出于可序列化的原因,我不得不使用java.util.Random
    【解决方案4】:

    使用 Spark Dataset API,可能用于累加器:

    df.withColumn("_n", substring(rand(),3,4).cast("bigint"))
    

    【讨论】:

      猜你喜欢
      • 1970-01-01
      • 1970-01-01
      • 2023-01-03
      • 2015-05-23
      • 2011-10-24
      • 1970-01-01
      • 1970-01-01
      • 1970-01-01
      相关资源
      最近更新 更多