sparkstreaming测试之三有状态的接收数据

测试思路:

网站建设哪家好,找成都创新互联!专注于网页设计、网站建设、微信开发、微信小程序定制开发、集团企业网站建设等服务项目。为回馈新老客户创新互联还提供了龙潭免费建站欢迎大家使用!

    首先,使用上篇文章的程序一发送网络数据;

    其次,运行spark程序,观察效果。

说明:

    1. 这里使用到了更新函数;

    2. 使用检查点来保证状态。

sparkStreaming

import org.apache.log4j.{LoggerLevel}
import org.apache.spark.streaming.{SecondsStreamingContext}
import org.apache.spark.{SparkContextSparkConf}
import org.apache.spark.streaming.StreamingContext._

object StatefulWordCount {
   def main(args:Array[]){

    Logger.().setLevel(Level.)
    Logger.().setLevel(Level.)

    updateFunc = (values: []state:Option[]) => {
      currentCount = values.foldLeft()(_+_)
      previousCount = state.getOrElse()
      (currentCount + previousCount)
    }

    conf = SparkConf().setAppName().setMaster()
    sc = SparkContext(conf)

    ssc = StreamingContext(sc())
    ssc.checkpoint()

    lines = ssc.socketTextStream(args()args().toInt)
    words = lines.flatMap(_.split())
    wordCounts = words.map(x=>(x))

    stateDstream = wordCounts.updateStateByKey[](updateFunc)
    stateDstream.print()
    ssc.start()
    ssc.awaitTermination()
  }
}

当前名称:sparkstreaming测试之三有状态的接收数据
本文URL:http://hxwzsj.com/article/jogpcs.html

其他资讯

Copyright © 2025 青羊区翔捷宏鑫字牌设计制作工作室(个体工商户) All Rights Reserved 蜀ICP备2025123194号-14
友情链接: 网站设计 成都网站制作 成都网站设计 H5网站制作 网站建设费用 网站建设公司 上市集团网站建设 营销型网站建设 重庆电商网站建设 成都网站制作 泸州网站建设 成都网站建设公司 专业网站建设 定制网站设计 成都网站建设 成都h5网站建设 企业网站设计 成都网站设计 网站建设改版 成都网站建设 成都商城网站建设 网站建设方案