尧图建网站 尧图建网站 YAOTU WEB BUILD 免费咨询
ARTICLE DETAIL

资讯详情

深耕网站建设与建站编程的一线实战洞察。

Spark大数据分析与实战笔记(第七章 Spark Streaming实时计算框架-03)

Spark大数据分析与实战笔记(第七章 Spark Streaming实时计算框架-03) 文章目录每日一句正能量第7章 Spark Streaming实时计算框架章节概要7.4 Spark Streaming整合Kafka实战7.4.1 KafkaUtils.createDstream方式7.4.2 KafkaUtils.createDirectStream方式每日一句正能量稳住自己的重心世界才会让你形成意义。世界的意义是一面镜子映照的是你自己的姿态。如果你的内心摇摆不定世界反射给你的便是混乱与虚无。只有当你内在稳定、拥有坚定的核心价值观、目标、自我认知你与世界的互动才会产生清晰的、属于你自己的“意义感”。第7章 Spark Streaming实时计算框架章节概要近年来在Web应用、网络监控、传感监测、电信金融、生产制造等领域增强了对数据实时处理的需求而Spark中的Spark Streaming实时计算框架就是为实现对数据实时处理的需求而设计。在电子商务中淘宝、京东网站从用户点击的行为和浏览的历史记录中发现用户的购买意图和兴趣然后通过Spark Streaming实时计算框架的分析处理为之推荐相关商品从而有效地提高商品的销售量同时也增加了用户的满意度可谓是“一举二得”。因此本章将针对Spark Streaming实时计算框架相关的知识进行详细介绍。7.4 Spark Streaming整合Kafka实战kafka作为一个实时的分布式消息队列实时地生产和消费消息。在这里我们可以利用Spark Streaming实时地读取Kafka中的数据然后再进行相关计算。在Spark1.3版本后KafkaUtils里面提供了两个创建DStream的方式一种是KafkaUtils.createDstream方式另一种为KafkaUtils.createDirectStream方式。本节我们针对DStream的这两种方式进行详细介绍。7.4.1 KafkaUtils.createDstream方式KafkaUtils.createDstream方式即基于Receiver的方式主要是通过Zookeeper连接Kafkareceivers接收器从Kafka中获取数据并且所有receivers获取到的数据都会保存在Spark executors中然后通过Spark Streaming启动job来处理这些数据具体处理流程如图7-12所示。图7-12在图中当Driver处理Spark Executors中的job时默认是会出现数据丢失的情况此时如果我们启用WAL日志将接收到数据同步地保存到分布式文件系统上如HDFS当数据由于某种原因丢失时丢失的数据能够及时恢复。接下来通过一个具体的案例来演示如何使用KafkaUtils.createDstream实现词频统计具体实现步骤如下导入依赖我们需要在pom.xml文件中添加Spark Streaming整合Kafka的依赖。具体内容如下dependencygroupIdorg.apache.spark/groupIdartifactIdspark-streaming-kafka-0-8_2.11/artifactIdversion2.0.2/version/dependency创建Scala类实现词频统计在spark_chapter07项目的/src/main/scala/cn.itcast.dstream目录下创建一个名为SparkStreaming_Kafka_createDstream的Scala类用来编写Spark Streaming应用程序实现词频统计。具体实现代码如文件7-1所示。文件7-1 SparkStreaming_Kafka_createDstream.scalaimportorg.apache.spark.streaming.dstream.{DStream,ReceiverInputDStream}importorg.apache.spark.streaming.kafka.KafkaUtilsimportorg.apache.spark.streaming.{Seconds,StreamingContext}importorg.apache.spark.{SparkConf,SparkContext}importscala.collection.immutableobjectSparkStreaming_Kafka_createDstream{defmain(args:Array[String]):Unit{//1.创建sparkConf,并开启wal预写日志保存数据源valsparkConf:SparkConfnewSparkConf().setAppName(SparkStreaming_Kafka_createDstream).setMaster(local[4]).set(spark.streaming.receiver.writeAheadLog.enable,true)//2.创建sparkContextvalscnewSparkContext(sparkConf)//3.设置日志级别sc.setLogLevel(WARN)//3.创建StreamingContextvalsscnewStreamingContext(sc,Seconds(5))//4.设置checkpointssc.checkpoint(./Kafka_Receiver)//5.定义zk地址valzkQuorumhadoop01:2181,hadoop02:2181,hadoop03:2181//6.定义消费者组valgroupIdspark_receiver//7.定义topic相关信息 Map[String, Int],这里的value并不是topic分区数,它表示的topic中每一个分区被N个线程消费valtopicsMap(kafka_spark-1)//8.通过KafkaUtils.createDstream对接kafka这时相当于同时开启3个receiver接受数据valreceiverDstream:immutable.IndexedSeq[ReceiverInputDStream[(String,String)]](1to3).map(x{valstream:ReceiverInputDStream[(String,String)]KafkaUtils.createStream(ssc,zkQuorum,groupId,topics)stream})//9.使用ssc中的union方法合并所有的receiver中的数据valunionDStream:DStream[(String,String)]ssc.union(receiverDstream)//10.SparkStreaming获取topic中的数据valtopicData:DStream[String]unionDStream.map(_._2)//11.按空格进行切分每一行,并将切分的单词出现次数记录为1valwordAndOne:DStream[(String,Int)]topicData.flatMap(_.split( )).map((_,1))//12.统计单词在全局中出现的次数valresult:DStream[(String,Int)]wordAndOne.reduceByKey(__)//13.打印输出结果result.print()//14.开启流式计算ssc.start()ssc.awaitTermination()}}运行代码后依次在hadoop01、hadoop02和hadoop03服务器执行命令zkServer.sh start启动Zookeeper集群然后依次在hadoop01、hadoop02和hadoop03服务器的Kafka根目录下执行命令bin/kafka-server-start.sh config/server.properties启动Kafka集群。结果如下图所示在使用Kafka发送消息和消费消息之前必须先要创建Topic用来指定消息的类别。具体命令如下kafka-topics.sh--create\--topickafka_spark\--partitions3\--replication-factor1\--zookeeperhadoop01:2181, hadoop02:2181,hadoop03:2181上述命令中创建了一个名为kafka_spark的Topic并且设置分区为3备份为1指定了Zookeeper集群的地址。结果如下图所示从上述内容可以看出名为kafka_spark的Topic已经创建完成。启动Kafka的消息生产者启动Kafka的消息生产者生产数据具体命令如下kafka-console-producer.sh\--broker-list hadoop01:9092\--topickafka_spark上述命令中我们指定消息生产者为hadoop01服务器。执行上述命令并且指定给kafka_spark这个Topic中发送消息。具体内容如下kafka-console-producer.sh\--broker-list hadoop01:9092\--topickafka_sparkhadoop spark hbase kafka sparkkafka itcast itcast spark kafka spark kafka结果如下图所示打开IDEA工具控制台输出的内容如图所示。注意如果我们使用KafkaUtils.createDstream方式时一开始系统会正常运行没有任何问题但是当系统出现异常重启SparkStreaming程序后则发现程序会重复处理已经处理过的数据。由于这种方式是使用Kafka的高级消费者APItopic的offset偏移量是在Zookeeper中。虽然这种方式会配合着WAL日志保证数据零丢失的高可靠性但却无法保证数据只被处理一次可能会处理两次。因此官方已经不推荐使用这种方式从而推荐我们使用KafkaUtils.createDirectStream方式。7.4.2 KafkaUtils.createDirectStream方式由于KafkaUtils.createDstream方式有一个弊端即无法保证数据只被处理一次因此我们来详细讲解官网推荐的方式即KafkaUtils.createDirectStream方式。KafkaUtils.createDirectStream方式不同于KafkaUtils.createDstream方式当接收数据时它会定期地从Kafka中Topic对应Partition中查询最新的偏移量再根据偏移量范围在每个batch里面处理数据然后Spark通过调用Kafka简单的消费者API即低级API来读取一定范围的数据具体处理流程如图所示。在图中当Driver处理Spark Executors中的job时系统突然出现异常重启Spark Streaming程序后程序会重复处理已经处理过的数据无法保证数据只被处理一次此时如果我们通过Spark中的StreamingContext对象将偏移量保存到CheckPoint中这样的话就可以避免因Spark Streaming和Zookeeper不同步即二者保存的偏移量不一致导致的数据被多次处理的现象。接下来通过一个具体的案例来演示如何使用KafkaUtils.createDirectStream方式来实现词频统计具体实现步骤如下导入依赖我们需要在pom.xml文件中添加Spark Streaming整合Kafka的依赖。同上这里不做赘述。具体内容如下dependencygroupIdorg.apache.spark/groupIdartifactIdspark-streaming-kafka-0-8_2.11/artifactIdversion2.0.2/version/dependency创建Scala类实现词频统计在spark_chapter07项目的/src/main/scala/cn.itcast.dstream目录下创建一个名为SparkStreaming_Kafka_createDirectStream的Scala类用来编写Spark Streaming应用程序实现词频统计。具体实现代码如文件7-8所示。文件7-8 SparkStreaming_Kafka_createDirectStream.scalaimportkafka.serializer.StringDecoderimportorg.apache.spark.streaming.dstream.{DStream,InputDStream}importorg.apache.spark.streaming.kafka.KafkaUtilsimportorg.apache.spark.streaming.{Seconds,StreamingContext}importorg.apache.spark.{SparkConf,SparkContext}objectSparkStreaming_Kafka_createDirectStream{defmain(args:Array[String]):Unit{//1.创建sparkConfvalsparkConf:SparkConfnewSparkConf().setAppName(SparkStreaming_Kafka_createDirectStream).setMaster(local[2])//2.创建sparkContextvalscnewSparkContext(sparkConf)//3.设置日志级别sc.setLogLevel(WARN)//4.创建StreamingContextvalsscnewStreamingContext(sc,Seconds(5))//5.设置chectPointssc.checkpoint(./Kafka_Direct)//6.配置kafka相关参数 metadata.broker.list为老版本的集群地址valkafkaParamsMap(metadata.broker.list-hadoop01:9092,hadoop02:9092,hadoop03:9092,group.id-spark_direct)//7.定义topicvaltopicsSet(kafka_direct0)//8.通过低级api方式将kafka与sparkStreaming进行整合valdstream:InputDStream[(String,String)]KafkaUtils.createDirectStream[String,String,StringDecoder,StringDecoder](ssc,kafkaParams,topics)//9.获取kafka中topic中的数据valtopicData:DStream[String]dstream.map(_._2)//10.按空格进行切分每一行,并将切分的单词出现次数记录为1valwordAndOne:DStream[(String,Int)]topicData.flatMap(_.split( )).map((_,1))//11.统计单词在全局中出现的次数valresult:DStream[(String,Int)]wordAndOne.reduceByKey(__)//12.打印输出结果result.print()//13.开启流式计算ssc.start()ssc.awaitTermination()}}运行代码后依次在hadoop01、hadoop02和hadoop03服务器执行命令zkServer.sh start启动Zookeeper集群然后依次在hadoop01、hadoop02和hadoop03服务器的Kafka根目录下执行命令bin/kafka-server-start.sh config/server.properties启动Kafka集群。创建Topic指定消息的类别在使用Kafka发送消息和消费消息之前必须先要创建Topic用来指定消息的类别。具体命令如下kafka-topics.sh--create\--topickafka_direct0\--partitions3\--replication-factor1\--zookeeperhadoop01:2181,hadoop02:2181,hadoop03:2181在上述命令中我们创建一个名为kafka_direct0的Topic并且设置分区为3备份为1指定了Zookeeper集群的地址。执行上述命令具体效果如下图所示从上述内容可以看出名为kafka_direct0的Topic已经创建完成。启动Kafka的消息生产者启动Kafka的消息生产者生产数据具体命令如下kafka-console-producer.sh\--broker-list hadoop01:9092\--topickafka_direct0上述命令中我们指定消息生产者为hadoop01服务器。执行上述命令并且指定给kafka_direct0这个Topic中发送消息即输入数据。具体内容如下kafka-console-producer.sh\--broker-list hadoop01:9092\--topickafka_direct0hadoop spark hbase kafka sparkkafka itcast itcast spark kafka spark kafka结果如下图所示打开IDEA工具控制台输出的内容如图7-15所示。使用KafkaUtils.createDirectStream方式的控制台输出。转载自https://blog.csdn.net/u014727709/article/details/163802560欢迎 点赞✍评论⭐收藏欢迎指正
返回列表