资讯

精准传达 • 有效沟通

从品牌网站建设到网络营销策划,从策略到执行的一站式服务

怎么解析SPARKforeach循环中的变量问题

怎么解析SPARK foreach循环中的变量问题,针对这个问题,这篇文章详细介绍了相对应的分析和解答,希望可以帮助更多想解决这个问题的小伙伴找到更简单易行的方法。

专业领域包括网站设计、网站制作商城网站建设、微信营销、系统平台开发, 与其他网站设计及系统开发公司不同,创新互联建站的整合解决方案结合了帮做网络品牌建设经验和互联网整合营销的理念,并将策略和执行紧密结合,为客户提供全网互联网整合方案。

原因

在spark算子中引用的外部变量,其实是变量的副本,在算子中对其值进行修改,只是改变副本的值,外部的变量还是没有变。
通俗易懂的讲就是foreach里的变量带不出来的,除非用map,将结果作为rdd返回

解决方案:

1、使用广播变量

object foreachtest {
  def main(args: Array[String]): Unit = {
 
    val conf = new SparkConf()
    conf.setMaster("local[1]")
    conf.setAppName("WcAppTask")
    val sc = new SparkContext(conf)
    sc.setLogLevel("WARN")
    val fileRdd = sc.parallelize(Array(("imsi1","2018-07-29 11:22:23","zd-A"),("imsi2","2018-07-29 11:22:24","zd-A"),("imsi3","2018-07-29 11:22:25","zd-A")))
    val result = mutable.Map.empty[String,String]
  val resultBroadCast: Broadcast[mutable.Map[String, String]] =sc.broadcast(result)
    fileRdd.foreach(input=>{
      val str = (input._1+"/t"+input._2+"/t"+input._3).toString
      resultBroadCast.value += (input._1.toString -> str)
      println(resultBroadCast.value.size) //返回1,2.3
    })
    println(result.size) //返回3
}

2:使用累加器

val accum = sc.collectionAccumulator[mutable.Map[String, String]]("My Accumulator")
fileRdd.foreach(input => {
  val str = input._1 + "/t" + input._2 + "/t" + input._3
  accum.add(mutable.Map(input._1 -> str))
})
println(accum.value.size())

3:累加变量 longAccumulator

val longaa= sc.longAccumulator("count")
fileRdd.foreach(input=>{
  val str = (input._1+"/t"+input._2+"/t"+input._3).toString
  longaa.add(1L)
})
println(longaa.count) //返回3

关于怎么解析SPARK foreach循环中的变量问题问题的解答就分享到这里了,希望以上内容可以对大家有一定的帮助,如果你还有很多疑惑没有解开,可以关注创新互联行业资讯频道了解更多相关知识。


网页题目:怎么解析SPARKforeach循环中的变量问题
文章源于:http://cdkjz.cn/article/pggshs.html
多年建站经验

多一份参考,总有益处

联系快上网,免费获得专属《策划方案》及报价

咨询相关问题或预约面谈,可以通过以下方式与我们联系

大客户专线   成都:13518219792   座机:028-86922220