第3步,运行分析器。在master节点开启终端,经由过程下面代码向Spark集群提交应用法度榜样。
- $ bin/spark-submit ~/StatefulWordCount.jar slave1 9999
查看结不雅
第1步,查看slave1上数据流模仿器运行情况。分析谱钥浏群上提走运行后与slave1上运行的数据流模仿器建立连接。当检测到外部连接时,数据流模仿器将每隔1000毫秒大年夜/home/dong/Streamingtext/file1.txt中随机朝长进步一行文本发送给slave1节点上的9999端口。因为该文本文件中每一行只包含一个悼?船是以每秒仅发送一个单词给端口。如图18所示。

图18 slave1上数据流模仿器运行示意图
第2步,查看master上分析器运行情况。在master节点的提交窗口中可以查看到统计结不雅,如图19所示。

图19 master上分析器运行示意图
图中注解截至147920770500ms分析器共接收到14个悼?船个中"spark"累计出现3次,"hbase"累计出现5次,"hello"累计出现3次,"world"累计出现3次。因为批处理时光距离是5s,模仿器每1秒发送1个悼?船使得分析器在5s内共接收到5个悼?船是以截止至147920771000ms,分析器共收到19个悼?船个中"spark"累计出现5次,"hbase"累计出现7次,"hello"累计出现4次,"world"累计出现3次。
第3步,查看HDFS中持久化目次。运行后查看HDFS上的持久化目次/user/dong/input/StatefulWordCountlog,如图20所示。Streaming应用法度榜样将接收到的收集数据持久化至该目次下,便于容错处理。

图20 HDFS上持久化目次示意图

window应用案例
在实际临盆情况中,与窗口相干的应用处景很常见,例如电商每距离10分钟小时统计某一商品前30分钟内累计发卖总额、趁魅站每隔1小时统计前3个小时内的客流量等,词攀类需求可借助Spark Streaming中的window相干操作实现。window应用案例同时涉及批处理时光距离、窗口时光距离与滑动时光距离。
功能需求
监听收集中某节点上指定端口传输的数据流(slave1节点上9999端口的英文文本数据,以逗号距离单词),每10秒统计前30秒各单词累计出现的次数。
代码实现
本例功能的实现涉及数据流模仿苹赝分析器两部分。
分析器代码:
- package dong.spark
- import org.apache.spark.{SparkContext, SparkConf}
- import org.apache.spark.streaming.StreamingContext._
- import org.apache.spark.streaming._
- import org.apache.spark.storage.StorageLevel
- object WindowWordCount {
- def main(args:Array[String]) ={
- val conf=new SparkConf().setAppName("WindowWordCount").
- setMaster("spark://192.168.149.132:7077")
- val sc=new SparkContext(conf)
- val ssc=new StreamingContext(sc, Seconds(5))
- ssc.checkpoint("hdfs://master:9000/user/dong/WindowWordCountlog")
- val lines=ssc.socketTextStream( args(0),
- args(1).toInt,
- StorageLevel.MEMORY_ONLY_SER)
- val words= lines.flatMap(_.split(","))
- /*采取reduceByKeyAndWindow操作进行叠加处理,窗口时光距离与滑动时光距离分别由参数args(2)和args(3)给出。*/
- val wordcounts=words.map(x=>(x,1)).
- reduceByKeyAndWindow((a:Int
推荐阅读
实现了功能虚拟化的收集可以或许使通信办事供给商快速供给办事、分析和主动化的收集,加快新办事投向市场的周期,并有效应用数据中间的通用平台。收集功能虚拟化旨在赞助电信行业加快立异>>>详细阅读
本文标题:大数据分析技术与实战之Spark Streaming
地址:http://www.17bianji.com/lsqh/37783.html
1/2 1

网友点评
精彩导读
科技快报
品牌展示