Spark大数据分析与实战笔记(第九章 综合案例—Spark实时交易数据统计-02)

Spark大数据分析与实战笔记(第九章 综合案例—Spark实时交易数据统计-02) 文章目录每日一句正能量第9章 综合案例—Spark实时交易数据统计章节概要9.3 模块开发—构建工程结构9.4 模块开发—构建订单系统9.4.1 模拟订单数据9.4.2 向Kafka集群发送订单数据9.5 模块开发 — 分析订单数据每日一句正能量活在自己的热爱里而不是别人的眼光里。热爱是自发燃烧的能量他人的评判常是扭曲的镜子。真正的自由始于将评价体系从外部收回手中。第9章 综合案例—Spark实时交易数据统计章节概要本章通过Spark Streaming技术开发商品实时交易数据统计模块案例该系统主要功能是在前端页面以动态报表展示后端不断增长的数据这也是所谓的看板平台。通过学习并开发看板平台从而帮助读者理解大数据实时计算架构的开发流程并能够掌握Spark实时计算框架Spark Streaming在实际应用中的使用方法。本章将 针对Spark实时交易数据统计进行详细讲解。9.3 模块开发—构建工程结构接下来我们分步骤讲解构建工程结构。创建工程首先打开IDEA开发工具创建Maven工程不选择任何模板具体如图9-5所示。图9-5然后单击【Next】按钮输入GroupId和ArtifactId作为组织名和项目工程名具体如图9-6所示。图9-6最后单击【Next】按钮直到出现【Finish】按钮完成工程创建。项目资源结构本项目中所涉及的包文件、配置文件以及页面文件等是项目中的组织结构如图9-7所示。图9-7我们将Spark工程和JavaWeb工程整合在一个Maven工程下因此还需要向项目中添加JavaWeb工程必备的web.xml文件。在IDEA开发工具中右键单击工程名选择Open Module Setting选项设置步骤如图9-8所示。图9-8在图中首先选择号添加Web模板然后依次修改路径和版本号并标记webapp路径最后点击【OK】按钮完成配置。添加依赖按照图9-7创建工程资源结构目录后在pom.xml配置文件中添加工程所需依赖具体代码如下所示。上述代码片段是项目所需的Spark依赖包含了spark-core、scala、spark-streaming和spark-streaming与kafka整合所需的jar文件。上述代码片段是项目所需Spring框架所需Jar文件。在上述代码片段是项目所需Jsp、Json数据转换工具、WebSocket的Jar文件。若读者仍需添加自己依赖库可通过https://mvnrepository.com/网站进行查找添加。9.4 模块开发—构建订单系统在本项目中我们利用Java编程构建订单系统在模拟订单数据时可以采用随机生成一组Json格式的字符串来模拟订单数据。9.4.1 模拟订单数据订单数据模型通常由订单编号、订单时间、商品编号、商品价格等数十个字段组成模型中的指标越多提供给分析人员可分析的维度就越多。首先在cn.itcast.createorder包下创建PaymentInfo.java文件用于定义订单字段以及生成订单数据具体代码如文件所示。文件9-1 PaymentInfo.javapackagecn.itcast.createorder;importcom.alibaba.fastjson.JSONObject;importjava.util.Random;importjava.util.UUID;publicclassPaymentInfo{privatestaticfinallongserialVersionUID1L;privateStringorderId;//订单编号privateStringproductId;//商品编号privatelongproductPrice;//商品价格publicPaymentInfo(){}publicstaticlonggetSerialVersionUID(){returnserialVersionUID;}publicStringgetOrderId(){returnorderId;}publicvoidsetOrderId(StringorderId){this.orderIdorderId;}publicStringgetProductId(){returnproductId;}publicvoidsetProductId(StringproductId){this.productIdproductId;}publiclonggetProductPrice(){returnproductPrice;}publicvoidsetProductPrice(longproductPrice){this.productPriceproductPrice;}OverridepublicStringtoString(){returnPaymentInfo{orderIdorderId\, productIdproductId\, productPriceproductPrice};}//模拟订单数据publicStringrandom(){RandomrnewRandom();this.orderIdUUID.randomUUID().toString().replaceAll(-,);this.productPricer.nextInt(1000);this.productIdr.nextInt(10);JSONObjectobjnewJSONObject();StringjsonStringobj.toJSONString(this);returnjsonString;}}模拟订单数据模块开发中地6-8行代码我们设置了三个字段分别是订单编号、商品编号、商品价格。第42-49行代码是模拟订单数据的核心方法我们采取使用UUID模拟生成订单编号UUID是由一组32位数的16进制数字随机构成的字符串数据商品编号是由0-9这十个数字组成代表特定商品。在数据传输过程中需要将对象转换成Json格式的字符串这里采用了Fastjson数据转换工具调用JSONObject类的toJSONString()方法将PaymentInfo订单对象转换为Json格式的字符串编写成功后就可以在test目录中创建测试用例最终随机生成的订单数据格式如下。orderId:b030e0dfb3b04cd18c3b32beac01ab25,productId:6,“productPrice:834}9.4.2 向Kafka集群发送订单数据模拟订单数据模块开发完成后接下来创建Kafka生产者对象将订单数据发送至Kafka集群中下面我们分步骤进行讲解。创建Kafka生产者对象在cn.itcast.createorder包下创建PaymentInfoProducer.java文件具体代码如文件9-2所示。文件9-2 PaymentInfoProducer.javapackagecn.itcast.createorder;importorg.apache.kafka.clients.producer.KafkaProducer;importorg.apache.kafka.clients.producer.ProducerRecord;importjava.util.Properties;publicclassPaymentInfoProducer{publicstaticvoidmain(String[]args){PropertiespropsnewProperties();// 1、指定Kafka集群的主机名和端口号props.put(bootstrap.servers,hadoop01:9092,hadoop02:9092,hadoop03:9092);// 2、指定等待所有副本节点的应答props.put(acks,all);// 3、指定消息发送最大尝试次数props.put(retries,0);// 4、指定一批消息处理大小props.put(batch.size,16384);// 5、指定请求延时props.put(linger.ms,1);// 6、指定缓存区内存大小props.put(buffer.memory,33554432);// 7、设置key序列化props.put(key.serializer,org.apache.kafka.common.serialization.StringSerializer);// 8、设置value序列化props.put(value.serializer,org.apache.kafka.common.serialization.StringSerializer);KafkaProducerString,StringkafkaProducernewKafkaProducerString,String(props);PaymentInfopaynewPaymentInfo();while(true){// 9、生产数据Stringmessagepay.random();kafkaProducer.send(newProducerRecordString,String(itcast_order,message));System.out.println(数据已发送到Kafakamessage);try{Thread.sleep(1000);}catch(InterruptedExceptione){e.printStackTrace();}}}}上述代码是利用Kafka API创建生产者对象设置Kafka集群配置参数并调用send()方法不断向指定Kafka集群中发送订单数据。2. 启动Kafka程序下面依次启动主机名为hadoop01、hadoop02、hadoop03这三台集群中的Kafka服务执行命令如下所示。bin/kafka-server-start.sh config/server.properties启动Kafka服务端进程后通过克隆hadoop01的会话窗口来创建名为itcast_order的Topic执行命令如下所示。kafka-topics.sh--create\--topicitcast_order\--partitions3\--replication-factor2\--zookeeperhadoop01:2181, hadoop02:2181, hadoop03:2181结果如下图所示Topic创建成功后就可以监听数据了执行命令如下所示。kafka-console-consumer.sh\--from-beginning--topicitcast_order\--bootstrap-server hadoop01:9092, hadoop02:9092, hadoop03:9092运行结果如下图所示命令执行完成后返回IDEA工具运行PaymentInfoProducer类生产数据随后观察Kafka消费数据的会话窗口和IDEA工具的控制台输出效果如图所示。9.5 模块开发 — 分析订单数据针对Kafka中的实时订单数据本节采用Spark Streaming实时计算框架对订单中不同商品的成交额进行统计分析然后将分析出的数据按照业务需求存入Redis数据库。配置Jedis操作Redis数据库数据写入到Redis可以使用Jedis工具Jedis是Redis官方推荐的Java连接开发工具其中集成了Redis操作命令、提供数据库的连接池管理以及使用简单等优点。在项目的资源目录创建redis.properties配置文件配置参数如文件9-3所示。文件9-3 redis.properties#表示jedis的服务器主机名jedis.hosthadoop01#表示jedis的服务的端口jedis.port6379#jedis连接池中最大的连接个数jedis.max.total60#jedis连接池中最大的空闲连接个数jedis.max.idle30#jedis连接池中最小的空闲连接个数jedis.min.idle5#jedis连接池最大的等待连接时间ms值jedis.max.wait.millis30000在scala目录的cn.itcast.processdata包下创建RedisClient.scala文件用于读取配置文件中Redis参数代码如文件9-4所示。文件9-4 RedisClient.scalapackagecn.itcast.processdataimportjava.util.Propertiesimportorg.apache.commons.pool2.impl.GenericObjectPoolConfigimportredis.clients.jedis.JedisPoolobjectRedisClient{valpropnewProperties()//加载配置文件prop.load(this.getClass.getClassLoader.getResourceAsStream(redis.properties))valredisHost:Stringprop.getProperty(jedis.host)valredisPort:Stringprop.getProperty(jedis.port)valredisTimeout:Stringprop.getProperty(jedis.max.wait.millis)lazyvalpoolnewJedisPool(newGenericObjectPoolConfig(),redisHost,redisPort.toInt,redisTimeout.toInt)lazyvalhooknewThread{overridedefrun{println(Execute hook thread: this)pool.destroy()}}}文件9-4是Scala版本的Jedis工具类为了读者掌握更多编程技巧同时提供了Java版本的Jedis工具类在cn.itcast.util包中创建JedisUtil.java文件用来操作Redis数据库具体代码如文件9-5所示。文件9-5 JedisUtil.javapackagecn.itcast.util;importredis.clients.jedis.Jedis;importredis.clients.jedis.JedisPool;importredis.clients.jedis.JedisPoolConfig;importjava.io.IOException;importjava.util.Properties;/** * Redis Java API 操作的工具类 * 主要为我们提供Java操作Redis的对象Jedis,类似数据库连接池 */publicclassJedisUtil{privateJedisUtil(){}privatestaticJedisPooljedisPool;static{PropertiespropnewProperties();try{prop.load(JedisUtil.class.getClassLoader().getResourceAsStream(redis.properties));JedisPoolConfigpoolConfignewJedisPoolConfig();//jedis连接池中最大的连接个数poolConfig.setMaxTotal(Integer.valueOf(prop.getProperty(jedis.max.total)));//jedis连接池中最大的空闲连接个数poolConfig.setMaxIdle(Integer.valueOf(prop.getProperty(jedis.max.idle)));//jedis连接池中最小的空闲连接个数poolConfig.setMinIdle(Integer.valueOf(prop.getProperty(jedis.min.idle)));//jedis连接池最大的等待连接时间ms值poolConfig.setMaxWaitMillis(Long.valueOf(prop.getProperty(jedis.max.wait.millis)));//表示jedis的服务器主机名Stringhostprop.getProperty(jedis.host);intportInteger.valueOf(prop.getProperty(jedis.port));jedisPoolnewJedisPool(poolConfig,host,port,10000);}catch(IOExceptione){e.printStackTrace();}}/** * 提供了Jedis的对象 * * return */publicstaticJedisgetJedis(){returnjedisPool.getResource();}/** * 资源释放 * * param jedis */publicstaticvoidreturnJedis(Jedisjedis){jedis.close();}}Spark Streaming处理数据接下来利用所学知识Spark Streaming处理Kafka集群中的数据在cn.itcast.processdata包下创建StreamintProcessdata.scala文件具体代码如文件9-6所示。文件9-6 StreamingProcessdata.scalapackagecn.itcast.processdataimportcom.alibaba.fastjson.{JSON,JSONObject}importkafka.serializer.StringDecoderimportorg.apache.spark.streaming.dstream.{DStream,InputDStream}importorg.apache.spark.streaming.kafka.KafkaUtilsimportorg.apache.spark.streaming.{Seconds,StreamingContext}importorg.apache.spark.{SparkConf,SparkContext}importredis.clients.jedis.JedisobjectStreamingProcessdata{//每件商品总销售额valorderTotalKeybussiness::order::total//总销售额valtotalKeybussiness::order::all//Redis数据库valdbIndex0defmain(args:Array[String]):Unit{//1、创建SparkConf对象valsparkConf:SparkConfnewSparkConf().setAppName(KafkaStreamingTest).setMaster(local[4])//2、创建SparkContext对象valscnewSparkContext(sparkConf)sc.setLogLevel(WARN)//3、构建StreamingContext对象valsscnewStreamingContext(sc,Seconds(3))//4、消息的偏移量就会被写入到checkpoint中ssc.checkpoint(./spark-receiver)//4、设置Kafka参数valkafkaParamsMap(bootstrap.servers-hadoop01:9092,hadoop02:9092,hadoop03:9092,group.id-spark-receiver)//5、指定Topic相关信息valtopicsSet(itcast_order)//6、通过KafkaUtils.createDirectStream利用低级api接受kafka数据valkafkaDstream:InputDStream[(String,String)]KafkaUtils.createDirectStream[String,String,StringDecoder,StringDecoder](ssc,kafkaParams,topics)//7、获取Kafka中Topic数据并解析JSON格式数据valevents:DStream[JSONObject]kafkaDstream.flatMap(lineSome(JSON.parseObject(line._2)))//按照productID进行分组统计个数和总价格valorders:DStream[(String,Int,Long)]events.map(x(x.getString(productId),x.getLong(productPrice))).groupByKey().map(x(x._1,x._2.size,x._2.reduceLeft(__)))orders.foreachRDD(xx.foreachPartition(partitionpartition.foreach(x{println(productIdx._1 countx._2 productPricricex._3)//获取Redis连接资源valjedis:JedisRedisClient.pool.getResource()//指定数据库jedis.select(dbIndex)//每个商品销售额累加jedis.hincrBy(orderTotalKey,x._1,x._3)//总销售额累加jedis.incrBy(totalKey,x._3)RedisClient.pool.returnResource(jedis)})))ssc.start()ssc.awaitTermination()}}上述代码中第16-26行代码用于构建StreamingContext对象并设置批处理时间间隔为3秒第27-36行代码设置Kafka连接参数并构建KafkaDstream对象通过KafkaUtils.createDirectStream()方法读取Kafka数据流第37-61行代码当接收到Kafka中每一条数据时通过JSON.parseObject()方法将Json字符串转换为JSONObject对象接着按照productId进行分组统计个数和价格将orders对象中的productId和productPrice字段以Hash数据类型的结构保存在Redis数据库中在Redis中表现为MaporderTotalKey, MapproductId, productPrice的数据格式。为了测试目前系统是否能够正常工作执行数据分析类StreamingProcessdata.scala、数据生产类PaymentInfoProducer最终在Redis客户端中查看数据如图9-10所示。图9-10 查看Redis数据从图9-10中可以看出数据成功保存在Redis数据库中。转载自https://blog.csdn.net/u014727709/article/details/163802748欢迎 点赞✍评论⭐收藏欢迎指正