发布者认证信息(营业执照和身份证)未完善,请登录后完善信息登录
 总算晓得SparkStreaming项目实战,实时计算Pv和Uv - 三农网
Hi,你好,欢迎来到三农网
  • 产品
  • 求购
  • 公司
  • 展会
  • 招商
  • 资讯
当前位置: 首页 » 资讯 » 重磅直击 找商家、找信息优选VIP,安全更可靠!
总算晓得SparkStreaming项目实战,实时计算Pv和Uv
发布日期:2021-11-13 13:30:35  浏览次数:4

最近有个需求,实时统计pv,uv,结果按照date,hour,pv,uv来展示,按天统计,第二天重新统计,当然了实际还需要按照类型字段分类统计pv,uv,比如按照date,hour,pv,uv,type来展示。这里介绍最基本的pv,uv的展示。

id uv pv date hour 1 2018-07-27 18

关于什么是pv,uv,可以参见这篇博客:/petermsh/article/details/

1、项目流程

日志数据从flume采集过来,落到hdfs供其它离线业务使用,也会sink到kafka,sparkStreaming从kafka拉数据过来,计算pv,uv,uv是用的redis的set集合去重,最后把结果写入mysql数据库,供前端展示使用。

2、具体过程 1)pv的计算

拉取数据有两种方式,基于received和direct方式,这里用direct直拉的方式,用的mapWithState算子保存状态,这个算子与updateStateByKey一样,并且性能更好。当然了实际中数据过来需要经过清洗,过滤,才能使用。

定义一个状态函数

// 实时流量状态更新函数   val mapFunction = (datehour:String, pv:Option[Long], state:State[Long]) => {     val accuSum = (0L) + ().getOrElse(0L)     val output = (datehour,accuSum)     (accuSum)     output   } 

这样就很容易的把pv计算出来了。

2)uv的计算

uv是要全天去重的,每次进来一个batch的数据,如果用原生的reduceByKey或者groupByKey对配置要求太高,在配置较低情况下,我们申请了一个93G的redis用来去重,原理是每进来一条数据,将date作为key,guid加入set集合,20秒刷新一次,也就是将set集合的尺寸取出来,更新一下数据库即可。

(rdd => {         (eachPartition => {         // 获取redis连接           val jedis = getJedis           (x => {             // 省略若干...             (key,)             // 设置存储每天的数据的set过期时间,防止超过redis容量,这样每天的set集合,定期会被自动删除             (key,)           })           // 关闭连接           closeJedis(jedis)         })       })  3)结果保存到数据库

结果保存到mysql,数据库,10秒刷新一次数据库,前端展示刷新一次,就会重新查询一次数据库,做到实时统计展示pv,uv的目的。

   def insertHelper(data: DStream[(String, Long)], tbName: String, colNames: String*): Unit = {     (rdd => {       val tmp_rdd = (x => (11, 13).toInt)       if (!()) {         val hour_now = () // 获取当前结果中最大的时间,在数据恢复中可以起作用         (eachPartition => {           try {             val jedis = getJedis             val conn = ()             (false)             val stmt = ()             (x => {               // val sql = ....               // 省略若干               (sql)             })             closeJedis(jedis)             () // 批量执行sql语句             ()             ()           } catch {             case e: Exception => {               (e)               ( + e)             }           }         })       }     })   }    // 计算当前时间距离次日零点的时长(毫秒) def resetTime = {     val now = new Date()     val todayEnd =      (, 23) //  12小时制     (, 59)     (, 59)     (, 999)      -   }  4)数据容错

流处理消费kafka都会考虑到数据丢失问题,一般可以保存到任何存储系统,包括mysql,hdfs,hbase,redis,zookeeper等到。这里用SparkStreaming自带的checkpoint机制来实现应用重启时数据恢复。

checkpoint

这里采用的是checkpoint机制,在重启或者失败后重启可以直接读取上次没有完成的任务,从kafka对应offset读取数据。

// 初始化配置文件 ()  val conf = new SparkConf().setAppName() ("","true") (".maxRatePerPartition",consumeRate) ("","24") val sc = new SparkContext(conf)  while (true){  val ssc = ( + (0),getStreamingContext _ )     ()     (resetTime)     (false,true) } 

checkpoint是每天一个目录,在第二天凌晨定时销毁StreamingContext对象,重新统计计算pv,uv。

注意:(false,true)表示优雅地销毁StreamingContext对象,不能销毁SparkContext对象,(true,true)会停掉SparkContext对象,程序就直接停了。

应用迁移或者程序升级

在这个过程中,我们把应用升级了一下,比如说某个功能写的不够完善,或者有逻辑错误,这时候都是需要修改代码,重新打jar包的,这时候如果把程序停了,新的应用还是会读取老的checkpoint,可能会有两个问题:

执行的还是上一次的程序,因为checkpoint里面也有序列化的代码;

直接执行失败,反序列化失败;

其实有时候,修改代码后不用删除checkpoint也是可以直接生效,经过很多测试,我发现如果对数据的过滤操作导致数据过滤逻辑改变,还有状态操作保存修改,也会导致重启失败,只有删除checkpoint才行,可是实际中一旦删除checkpoint,就会导致上一次未完成的任务和消费kafka的offset丢失,直接导致数据丢失,这种情况下我一般这么做。

这种情况一般是在另外一个集群,或者把checkpoint目录修改下,我们是代码与配置文件分离,所以修改配置文件checkpoint的位置还是很方便的。然后两个程序一起跑,除了checkpoint目录不一样,会重新建,都插入同一个数据库,跑一段时间后,把旧的程序停掉就好。以前看官网这么说,只能记住不能清楚明了,只有自己做时才会想一下办法去保证数据准确。

5)保存offset到mysql

如果保存offset到mysql,就可以将pv, uv和offset作为一条语句保存到mysql,从而可以保证exactly-once语义。

var messages: InputDStream[ConsumerRecord[String, String]] = null       if () {         messages = [String, String](           ssc           ,            , [String, String](topics, kafkaParams, )         )       } else {          messages = [String, String](           ssc           ,            , [String, String](topics, kafkaParams)         )       }               (rdd => {           .... }) 

从mysql读取offset并且解析:

   def getLastOffsets(tbName: String): [TopicPartition, Long] = {     val sql = s"select offset from ${tbName} where id = (select max(id) from ${tbName})"     val conn = (config)     val psts = (sql)     val res = ()     var tpMap: [TopicPartition, Long] = [TopicPartition, Long]()     while (()) {       val o = (1)       val jSONArray = (o)       ().foreach(offset => {         val json = ()         val topicAndPartition = new TopicPartition(("topic"), ("partition"))         (topicAndPartition, ("untilOffset"))       })     }     (res, psts, conn)     tpMap }  6)日志

日志用的log4j2,本地保存一份,ERROR级别的日志会通过邮件发送到手机,如果错误太多也会被邮件轰炸,需要注意。

val logger = ()   // 邮件level=error日志   val logger2 = ("email") 

 

VIP企业最新发布
全站最新发布
最新VIP企业
背景开启

三农网是一个开放的平台,信息全部为用户自行注册发布!并不代表本网赞同其观点或证实其内容的真实性,需用户自行承担信息的真实性,图片及其他资源的版权责任! 本站不承担此类作品侵权行为的直接责任及连带责任。

如若本网有任何内容侵犯您的权益,请联系 QQ: 1130861724

网站首页 | 实时热点 | 侵权删除 | 付款方式 | 联系方式 | 法律责任 | 网站地图 ©2022 zxb2b.com 三农网,中国大型农产品交易电商平台 鄂公网安备42018502006996 SITEMAPS | 鄂ICP备14015623号-20

返回顶部