代码之家  ›  专栏  ›  技术社区  ›  lu_ferra

为Spark Streaming Scala分配变量值

  •  1
  • lu_ferra  · 技术社区  · 9 年前

    我必须计算数据流中元组的数量,并根据值修改布尔变量的值。 不幸的是,我所做的似乎不是分配新的值。

    val teon = false
    
    s1.foreachRDD( rdd => {
    System.out.println("# events = " + rdd.count())
      if (rdd.count().>(1000)) 
        teon.equals(true) 
      else 
        teon.equals(false)
    })
    
    if(teon){
     val ton2 = s2.map { x => x.sensor_name }
     ton2.print
    }
    else {
      val ton3 = s2.map { x => x.stt.spatial.unit }
      ton3.print
    }
    

    s1 s2 是数据流[传感器](传感器是自定义类)。

    2 回复  |  直到 9 年前
        1
  •  1
  •   maasg    9 年前

    该代码中有两个问题:

    第一个是 teon 声明为 val 。它是不可变的,因此,它的值在程序执行过程中永远不会改变。

    DStream

    if(teon){
     val ton2 = s2.map { x => x.sensor_name }
     ton2.print
    }
    

    当程序首次加载并添加到数据流转换DAG以执行时,将仅评估一次。让我们记住 数据流 编程模型基于当 streamingContext 这一点以后不会改变。

    因为我们想根据价值观做出动态选择 包含 在这个过程中,我们需要做出这些决定 数据流 活动

    var teon = false
    
    s1.foreachRDD{ rdd => 
      val count = rdd.count // compute it only once!
      System.out.println("# events = " + count)
      teon  =  count > 1000 // use the boolean value directly
    }
    
    s2.foreachRDD { rdd =>  
      val ton = if (teon) { 
        rdd.map( x => x.sensor_name )
      } else {
        rdd.map( x => x.stt.spatial.unit ) // I'm assuming here that sensor_name and _stt.spatial_unit are the same type.
      }
      ton.take(10).foreach(e => println(e)) // implement DStream.print "by hand"
    }
    
        2
  •  1
  •   Yehor Krivokon    9 年前
    val teon = false    
    ...   
    ...    
    ...    
    if(teon)   
    

    这没有道理。
    增值税 意味着你不能改变变量的值(就像 最终的 在Java中)。所以它总是假的。

     var teon = false
    
    推荐文章