news 2026/8/25 15:54:59

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

作者头像

张小明

前端开发工程师

1.2k 24
文章封面图
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 模块开发—构建工程结构

接下来,我们分步骤讲解构建工程结构。

  1. 创建工程
    首先打开IDEA开发工具,创建Maven工程,不选择任何模板,具体如图9-5所示。

图9-5

然后单击【Next】按钮,输入GroupId和ArtifactId,作为组织名和项目工程名,具体如图9-6所示。

图9-6

最后单击【Next】按钮直到出现【Finish】按钮完成工程创建。

  1. 项目资源结构
    本项目中所涉及的包文件、配置文件以及页面文件等是项目中的组织结构,如图9-7所示。

图9-7

我们将Spark工程和JavaWeb工程整合在一个Maven工程下,因此还需要向项目中添加JavaWeb工程必备的web.xml文件。在IDEA开发工具中,右键单击工程名,选择"Open Module Setting"选项,设置步骤如图9-8所示。

图9-8

在图中,首先选择"+"号添加Web模板,然后依次修改路径和版本号,并标记webapp路径,最后点击【OK】按钮完成配置。

  1. 添加依赖
    按照图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.java

packagecn.itcast.createorder;importcom.alibaba.fastjson.JSONObject;importjava.util.Random;importjava.util.UUID;publicclassPaymentInfo{privatestaticfinallongserialVersionUID=1L;privateStringorderId;//订单编号privateStringproductId;//商品编号privatelongproductPrice;//商品价格publicPaymentInfo(){}publicstaticlonggetSerialVersionUID(){returnserialVersionUID;}publicStringgetOrderId(){returnorderId;}publicvoidsetOrderId(StringorderId){this.orderId=orderId;}publicStringgetProductId(){returnproductId;}publicvoidsetProductId(StringproductId){this.productId=productId;}publiclonggetProductPrice(){returnproductPrice;}publicvoidsetProductPrice(longproductPrice){this.productPrice=productPrice;}@OverridepublicStringtoString(){return"PaymentInfo{"+"orderId='"+orderId+'\''+", productId='"+productId+'\''+", productPrice="+productPrice+'}';}//模拟订单数据publicStringrandom(){Randomr=newRandom();this.orderId=UUID.randomUUID().toString().replaceAll("-","");this.productPrice=r.nextInt(1000);this.productId=r.nextInt(10)+"";JSONObjectobj=newJSONObject();StringjsonString=obj.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集群中,下面我们分步骤进行讲解。

  1. 创建Kafka生产者对象
    在cn.itcast.createorder包下创建PaymentInfoProducer.java文件,具体代码如文件9-2所示。

文件9-2 PaymentInfoProducer.java

packagecn.itcast.createorder;importorg.apache.kafka.clients.producer.KafkaProducer;importorg.apache.kafka.clients.producer.ProducerRecord;importjava.util.Properties;publicclassPaymentInfoProducer{publicstaticvoidmain(String[]args){Propertiesprops=newProperties();// 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");KafkaProducer<String,String>kafkaProducer=newKafkaProducer<String,String>(props);PaymentInfopay=newPaymentInfo();while(true){// 9、生产数据Stringmessage=pay.random();kafkaProducer.send(newProducerRecord<String,String>("itcast_order",message));System.out.println("数据已发送到Kafaka:"+message);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数据库。

  1. 配置Jedis操作Redis数据库
    数据写入到Redis,可以使用Jedis工具,Jedis是Redis官方推荐的Java连接开发工具,其中集成了Redis操作命令、提供数据库的连接池管理以及使用简单等优点。

在项目的资源目录创建redis.properties配置文件,配置参数如文件9-3所示。

文件9-3 redis.properties

#表示jedis的服务器主机名jedis.host=hadoop01#表示jedis的服务的端口jedis.port=6379#jedis连接池中最大的连接个数jedis.max.total=60#jedis连接池中最大的空闲连接个数jedis.max.idle=30#jedis连接池中最小的空闲连接个数jedis.min.idle=5#jedis连接池最大的等待连接时间ms值jedis.max.wait.millis=30000

在scala目录的cn.itcast.processdata包下创建RedisClient.scala文件,用于读取配置文件中Redis参数,代码如文件9-4所示。

文件9-4 RedisClient.scala

packagecn.itcast.processdataimportjava.util.Propertiesimportorg.apache.commons.pool2.impl.GenericObjectPoolConfigimportredis.clients.jedis.JedisPoolobjectRedisClient{valprop=newProperties()//加载配置文件prop.load(this.getClass.getClassLoader.getResourceAsStream("redis.properties"))valredisHost:String=prop.getProperty("jedis.host")valredisPort:String=prop.getProperty("jedis.port")valredisTimeout:String=prop.getProperty("jedis.max.wait.millis")lazyvalpool=newJedisPool(newGenericObjectPoolConfig(),redisHost,redisPort.toInt,redisTimeout.toInt)lazyvalhook=newThread{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.java

packagecn.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{Propertiesprop=newProperties();try{prop.load(JedisUtil.class.getClassLoader().getResourceAsStream("redis.properties"));JedisPoolConfigpoolConfig=newJedisPoolConfig();//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的服务器主机名Stringhost=prop.getProperty("jedis.host");intport=Integer.valueOf(prop.getProperty("jedis.port"));jedisPool=newJedisPool(poolConfig,host,port,10000);}catch(IOExceptione){e.printStackTrace();}}/** * 提供了Jedis的对象 * * @return */publicstaticJedisgetJedis(){returnjedisPool.getResource();}/** * 资源释放 * * @param jedis */publicstaticvoidreturnJedis(Jedisjedis){jedis.close();}}
  1. Spark Streaming处理数据
    接下来利用所学知识Spark Streaming,处理Kafka集群中的数据,在cn.itcast.processdata包下创建StreamintProcessdata.scala文件,具体代码如文件9-6所示。

文件9-6 StreamingProcessdata.scala

packagecn.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{//每件商品总销售额valorderTotalKey="bussiness::order::total"//总销售额valtotalKey="bussiness::order::all"//Redis数据库valdbIndex=0defmain(args:Array[String]):Unit={//1、创建SparkConf对象valsparkConf:SparkConf=newSparkConf().setAppName("KafkaStreamingTest").setMaster("local[4]")//2、创建SparkContext对象valsc=newSparkContext(sparkConf)sc.setLogLevel("WARN")//3、构建StreamingContext对象valssc=newStreamingContext(sc,Seconds(3))//4、消息的偏移量就会被写入到checkpoint中ssc.checkpoint("./spark-receiver")//4、设置Kafka参数valkafkaParams=Map("bootstrap.servers"->"hadoop01:9092,hadoop02:9092,hadoop03:9092","group.id"->"spark-receiver")//5、指定Topic相关信息valtopics=Set("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(line=>Some(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(x=>x.foreachPartition(partition=>partition.foreach(x=>{println("productId="+x._1+" count="+x._2+" productPricrice="+x._3)//获取Redis连接资源valjedis:Jedis=RedisClient.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中表现为Map<orderTotalKey, Map<productId, productPrice>>的数据格式。

为了测试目前系统是否能够正常工作,执行数据分析类(StreamingProcessdata.scala)、数据生产类(PaymentInfoProducer),最终在Redis客户端中查看数据如图9-10所示。

图9-10 查看Redis数据

从图9-10中可以看出,数据成功保存在Redis数据库中。


转载自:https://blog.csdn.net/u014727709/article/details/163802748
欢迎 👍点赞✍评论⭐收藏,欢迎指正

版权声明: 本文来自互联网用户投稿,该文观点仅代表作者本人,不代表本站立场。本站仅提供信息存储空间服务,不拥有所有权,不承担相关法律责任。如若内容造成侵权/违法违规/事实不符,请联系邮箱:809451989@qq.com进行投诉反馈,一经查实,立即删除!
网站建设 2026/8/25 15:44:37

rocketMQ proxy 延迟队列

Proxy 本身不做延迟队列的存储和调度 &#xff08;那是 Broker 端 ScheduleMessageService / TimerMessageStore 的职责&#xff09;&#xff0c;Proxy 只负责 延迟消息的「发送端属性填充 延迟级别换算」以及消费端识别 。相关实现集中在 gRPC 发送链路和配置里。 发送端&…

作者头像 李华
网站建设 2026/8/25 15:24:30

机器人破人类百米纪录:四足机器人的运动控制与工业落地

2026年&#xff0c;机器人再次刷新公众对运动能力的认知&#xff1a;一台四足机器人以10米/秒级别的速度完成百米冲刺&#xff0c;将“机器人破人类百米纪录”从实验室话题推向工程现实。对政企采购决策者而言&#xff0c;真正值得关注的不是单次赛跑成绩&#xff0c;而是这类高…

作者头像 李华