Streaming与kafkaupdateStateBykey()

object H extends App{
        val  conf=new  SparkConf().setMaster("local[2]").setAppName("hello")
        val ss=new StreamingContext(conf,Seconds(5))
        val kafkaParams=Map[String,String]("metadata.broker.list"->"myhadoop1:9092")
        ss.checkpoint("hdfs://myhadoop1:8020/data")
        val topic=Set[String]("wordcount1")
        //kafka
        val lines=KafkaUtils.createDirectStream[String,String,StringDecoder,StringDecoder](ss,kafkaParams,topic)
        lines.flatMap(_._2.split(" ")).map((_,1)).updateStateByKey((seqs:Seq[Int],option:Option[Int])=>{
                var oldValue=option.getOrElse(0)
                for(seq<-seqs){
                        oldValue+=seq
                }
                Option[Int](oldValue)
        }).print()
        ss.start()
        ss.awaitTermination()
}

本文名称:Streaming与kafkaupdateStateBykey()
URL标题:http://www.hxwzsj.com/article/gcojeh.html

其他资讯

Copyright © 2025 青羊区翔捷宏鑫字牌设计制作工作室(个体工商户) All Rights Reserved 蜀ICP备2025123194号-14
友情链接: 成都网站建设 高端定制网站设计 成都网站设计 高端网站设计推广 企业网站设计 网站制作 定制网站建设 成都网站建设 成都网站制作 重庆企业网站建设 成都网站设计制作公司 教育网站设计方案 盐亭网站设计 成都网站设计公司 手机网站制作 重庆网站建设 成都响应式网站建设 网站建设公司 成都网站设计 企业网站建设 定制网站设计 成都企业网站设计