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

发布时间:2026/8/25 15:55:23

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欢迎 点赞✍评论⭐收藏欢迎指正

相关新闻

Python 详解:从语法基础到进阶实战

Python 详解:从语法基础到进阶实战

2026/8/25 15:55:23

1. Python 简介 Python 是一门简洁、易读、功能强大的高级编程语言,由 Guido van Rossum 在 1991 年首次发布。它强调代码可读性,用缩进表达代码块,拥有丰富的标准库和第三方生态,被广泛应用于 Web 开发、数据分析、人工智能、自动…

PyTorch深度学习与实践【03】【数据的三种类型及其编码方案】

PyTorch深度学习与实践【03】【数据的三种类型及其编码方案】

2026/8/25 15:55:23

一、数据的三种类型 (一)连续值(比例 / 区间尺度) 数值之间的差值、倍数有实际物理含义。例子:重量 3kg,10kg。10‑37,代表重量差 7kg;10kg 是 3kg 的三倍重。 葡萄酒里面酒精度、酸…

辽宁智慧校园平台建设方案怎么选?几点实用经验帮你少走弯路

辽宁智慧校园平台建设方案怎么选?几点实用经验帮你少走弯路

2026/8/25 15:45:22

✅作者简介:合肥自友科技 📌核心产品:智慧校园平台(包括教工管理、学工管理、教务管理、考务管理、后勤管理、德育管理、资产管理、公寓管理、实习管理、就业管理、离校管理、科研平台、档案管理、学生平台等26个子平台) 。公司所有人员均有多…

【非标自动化】3、AutoShop快速理解(系统变量表)

【非标自动化】3、AutoShop快速理解(系统变量表)

2026/8/25 16:35:25

这些“系统变量表”不只是给你看的,它们本质上是 AutoShop 已经预先定义好的一批特殊变量,用来让程序访问 PLC 自身的状态、通信状态、模块状态和系统参数。可以先把变量分成两类:普通变量:由用户自己定义 系统变量:由…

【非标自动化】3、AutoShop快速理解(软件界面)

【非标自动化】3、AutoShop快速理解(软件界面)

2026/8/25 16:35:25

这张界面可以先不要把它看成“很多复杂菜单”,而要把它理解成一套完整的PLC工程工作台:先配置PLC和硬件↓ 定义变量和设备地址↓ 编写控制程序↓ 编译检查↓ 连接PLC并下载↓ 在线监控和调试↓ 排查故障、保存项目界面上的不同区域,就是分别服…

黑马苍穹外卖笔记day10

黑马苍穹外卖笔记day10

2026/8/25 16:35:25

订单状态定时处理、来单提醒和客户催单Spring Task:Spring Task是Spring框架提供的任务调度工具,可以按照约定的时间自动执行某个代码逻辑。应用场景特别广泛,只要是需要定制处理的场景都可以使用Spring Taskcron表达式:cron表达式…

【非标自动化】3、AutoShop快速理解(软件介绍)

【非标自动化】3、AutoShop快速理解(软件介绍)

2026/8/25 16:35:25

AutoShop是汇川面向小型PLC(Programmable Logic Controller,可编程逻辑控制 器)产品的编程组态软件,具有友好的编程和调试环境,拥有丰富、强大的通信和控 制功能。支持梯形图(LD)、顺序功能图&a…

Kotlin 语言【知识点整理2】

Kotlin 语言【知识点整理2】

2026/8/25 16:35:24

目录 一、基本概念 1.包的定义与导入 2.程序入口点 2.1 输入 3.变量 二、基本类型 1.数字 1.1 整数类型 1.2 浮点类型 1.3 数字字面常量 1.4 装箱与缓存 1.4.1 JVM是怎么存储数字的? 1.4.2 使用可空类型的时候会触发装箱操作: 1.4.3 JVM 对…

es怎么做拆词的

es怎么做拆词的

2026/8/25 16:25:24

ES 做拆词(分词)的核心是分词器,它在索引文档和搜索时,负责把长文本切成一个个独立的词(term),这样才能建立倒排索引,实现高效的全文搜索。 这个过程主要有三种角色: 输入…

[光学原理与应用-521]:对光的错误理解与纠偏

[光学原理与应用-521]:对光的错误理解与纠偏

2026/8/24 19:53:32

首先光是一种能量的载体和形态,宏观上观察到的光是由无数个微观的光量子组成的,每个光子在产生的瞬间,其在真空的空间中以确定不变的速度沿着一个初始的方向一直向前,在微观层面,每个光量子的运动轨迹是以波函数所展现…

SIP通话转接原理与REFER方法实战解析

SIP通话转接原理与REFER方法实战解析

2026/8/24 19:56:07

1. 通话转接不是“挂断再拨号”,而是SIP会话的动态重定向你有没有遇到过这样的场景:客服坐席A正在和客户通电话,突然需要把这通对话无缝转给专家坐席B,客户完全感知不到中间的断连——既没听到忙音,也没被要求重新拨号…

Kolla-ansible单节点OpenStack部署实战:从环境准备到排坑指南

Kolla-ansible单节点OpenStack部署实战:从环境准备到排坑指南

2026/8/24 21:16:09

1. 为什么选择Kolla-ansible来部署单节点OpenStack?如果你正在寻找一种能把OpenStack从“概念”快速变成“可用的实验环境”的方法,那么Kolla-ansible几乎是当前最主流、最省心的选择。我见过太多人卡在手动编译依赖、配置服务、处理版本冲突的泥潭里&am…

三步把QQ空间历史说说导出到本地:GetQzonehistory 极简指南

三步把QQ空间历史说说导出到本地:GetQzonehistory 极简指南

2026/8/25 0:04:34

三步把QQ空间历史说说导出到本地:GetQzonehistory 极简指南 【免费下载链接】GetQzonehistory 获取QQ空间发布的历史说说 项目地址: https://gitcode.com/GitHub_Trending/ge/GetQzonehistory Meta Description:GetQzonehistory 是一个QQ空间历史说…

洛谷 P7912:[CSP-J 2021 T4] 小熊的果篮 ← 双向链表

洛谷 P7912:[CSP-J 2021 T4] 小熊的果篮 ← 双向链表

2026/8/25 0:04:35

【题目来源】 https://www.luogu.com.cn/problem/P7912 【题目描述】 小熊的水果店里摆放着一排 n 个水果。每个水果只可能是苹果或桔子,从左到右依次用正整数 1,2,…,n 编号。连续排在一起的同一种水果称为一个“块”。小熊要把这一排水果挑到若干个果篮里&#x…

Transformers.js 网页端图像抠图实战:零后端 3 行代码返回透明 PNG

Transformers.js 网页端图像抠图实战:零后端 3 行代码返回透明 PNG

2026/8/25 0:04:35

Transformers.js 网页端图像抠图实战:零后端 3 行代码返回透明 PNG 【免费下载链接】transformers.js State-of-the-art Machine Learning for the web. Run 🤗 Transformers directly in your browser, with no need for a server! 项目地址: https:/…

摆脱论文困扰!盘点2026年全网爆红的的AI论文写作工具

摆脱论文困扰!盘点2026年全网爆红的的AI论文写作工具

2026/8/22 2:02:26

一天写完毕业论文在2026年已不再是天方夜谭。2026年最炸裂、实测能大幅提速的AI论文写作工具,覆盖选题构思、文献整理、内容生成、格式排版等核心场景,真正帮你高效搞定论文难题。 一、全流程王者:一站式搞定论文全链路(一天定稿首…

导师推荐!2026最新AI论文工具测评与实用推荐

导师推荐!2026最新AI论文工具测评与实用推荐

2026/8/22 4:13:47

2026年真正好用的AI论文工具,核心看生成的论文质量、低AI味、格式正确、学术适配四大指标。综合实测,千笔AI、ThouPen、豆包、DeepSeek、Grammarly 是当前最值得推荐的梯队,覆盖从免费到付费、从中文到英文、从文科到理工的全场景需求。 一、…

告别游戏崩溃:XCOM 2模组管理器的智能革命

告别游戏崩溃:XCOM 2模组管理器的智能革命

2026/8/22 1:32:34

告别游戏崩溃:XCOM 2模组管理器的智能革命 【免费下载链接】xcom2-launcher The Alternative Mod Launcher (AML) is a replacement for the default game launchers from XCOM 2 and XCOM Chimera Squad. 项目地址: https://gitcode.com/gh_mirrors/xc/xcom2-lau…