模块创建和数据准备
在UserBehaviorAnalysis下,构建了一个新的maven模块作为子项目,它的名字是NetworkFlowAnalysis。在这个子模块中,我们也没有引入更多的依赖项,所以我们不需要更改pom文件。

在src/main/ directory下,将默认的源文件目录java重命名为scala。将apache服务器的日志文件apache.log复制到资源文件目录src/main/resources,我们将从这里读取数据。
当然,我们仍然可以使用UserBehavior.csv作为数据源。这个时候我们不会分析每一个对服务器的访问请求,而是具体的页面浏览操作。
基于服务器日志的热门页面访问量统计
我们现在要实现的模块是“实时流量统计”。对于一个电商平台来说,用户登录的入口流量和不同页面的访问流量都是值得分析的重要数据,这些数据可以简单的从web服务器的日志中提取出来。
这里我们先实现“热门页面浏览量”的统计,即读取服务器日志中的每条日志,统计一段时间内用户访问每个url的次数,然后排序输出显示。
具体方法是:每隔5秒,输出最近10分钟内访问量最大的前N个网址。可以看出,这个要求和前面的“实时热门商品统计”很像,可以借鉴前面的代码。
在src/main/scala下创建NetworkFlow.scala文件,并创建一个新的singleton对象。定义样本类ApacheLogEvent,它是输入日志数据流;还有UrlViewCount,是窗口操作统计的输出数据类型。在main函数中创建StreamExecutionEnvironment并进行配置,然后从apache.log文件中读取数据,包装成ApacheLogEvent类型。
需要注意的是,原始日志中的时间是“dd/MM/yyyy:HH:mm:ss”的形式,需要定义一个DateTimeFormat将其转换成我们需要的时间戳格式:
。地图
val sdf =新的简单日期格式
val timestamp = sdf.parse)。取得时间
ApacheLogEvent,linearray,timestamp,
线性阵列,线性阵列)
})
完整的代码如下:
networkflow analysis/src/main/Scala/network flow . Scala
案例类ApacheLogEvent
案例类UrlViewCount
对象网络流{
def main: Unit = {
val env = streamexecutionenvironment . getexecutionenvironment
env . setstreamtimecharacter istic
env.setParallelism
val流=环境
//以窗口为例,用自己的路径替换。
。readTextFile
。地图
val simple date format = new simple date format
val timestamp = simple date format . parse)。取得时间
ApacheLogEvent,linearray,timestamp,linearray,
线性阵列)
})
。assignitemestampsandwatermark){
覆盖定义提取时间戳:Long = {
t .事件时间
}
})
。过滤器$)。)*$".r
。非空的
} )
。基比
。时间窗口,时间.秒)
。聚合,新WindowResultFunction)
。基比
。流程)
。打印
环境执行
}
类CountAgg扩展aggregate function[Apache logevent,Long,Long] {
覆盖def createAccumulator: Long = 0L
覆盖定义添加:Long = acc + 1
覆盖def getResult: Long = acc
覆盖定义合并:Long = acc1 + acc2
}
类WindowResultFunction扩展WindowFunction[Long,UrlViewCount,Tuple,
时间窗口] {
覆盖定义应用:单位= {
val URL:String = key . as instance of[tuple 1[String]]. F0
val count = aggregate result . iterator . next
收集器.收集)
}
}
类TopNHotUrls扩展了KeyedProcessFunction[Tuple,UrlViewCount,
字符串] {
private var urlState:ListState[UrlViewCount]= _
覆盖定义打开:单位= {
超级开放
val urlStateDesc = new ListStateDescriptor[UrlViewCount]
URL state = getruntimecontext . getliststate
}
覆盖定义过程元素:单位= {
//每条数据都保存到状态
urlState.add
context . timer service . registerevent timer
}
覆盖def onTimer:单位= {
//获取收到的所有URL访问
val allUrlViews:list buffer[UrlViewCount]= list buffer
导入Scala . collection . Java conversions . _
对于{
allUrlViews += urlView
}
//提前清除状态中的数据,释放空
urlState.clear
//按访问次数从大到小排序
val sorted urlviews = allurlviews . sort by
。拿
//将排名信息格式化为字符串,以便于打印
var结果:StringBuilder = new StringBuilder
结果.追加
result.append.append)。附加
对于{
val current urlview:UrlViewCount = sorted urlviews
//例如no1:URL =/blog/tags/Firefox flav = RSS 20 traffic = 55
结果.追加.追加.追加
.追加.追加
追加追加追加
}
结果.追加
//控制输出频率,模拟实时滚动结果。
线程.睡眠
out.collect
}
}
}
基于隐藏日志数据的网络流量统计
我们发现从web服务器日志中获取的url经常会要求一个资源地址,如果要统计页数的话往往需要过滤。在实际的电子商务应用中,我们可能更关心整个电子商务网站的网络流量,而不是各个单独页面的流量。
这个指标,除了合并前每页的统计结果外,还可以通过统计埋藏日志数据中的“pv”行为得到。
网站总浏览量统计
衡量网站流量最简单的指标之一就是网站的浏览量。用户每打开一个页面,就记录一次PV,多次打开同一个页面,浏览量就累积起来了。一般来说,PV与访客数量成正比,但PV并不直接决定页面的真实访客数量。就像一个访问者可以通过不断刷新页面来创造非常高的PV一样。
我们知道,当用户浏览一个页面时,他会从浏览器向web服务器发送一个请求。在接收到这个请求后,web服务器将把对应于这个请求的网页发送给浏览器,从而产生一个PV。所以我们的统计方法可以是从web服务器的日志中提取相应的页面访问量然后进行统计,就像上一节一样;也可以直接从掩埋日志中提取用户发送的页面请求,从而统计总浏览量。
所以,接下来,我们使用UserBehavior.csv作为数据源,实现一个网站总浏览量的统计。我们可以设置一个滚动的时间窗口,实时统计网站的PV。
在src/main/scala下创建PageView.scala文件,具体代码如下:
networkflow analysis/src/main/Scala/pageview . Scala
案例类用户行为
对象页面视图{
def main: Unit = {
val resources path = getclass . get resource
val env = streamexecutionenvironment . getexecutionenvironment
env . setstreamtimecharacter istic
env.setParallelism
val stream = env.readTextFile
。地图
UserBehavior.toLong,dataArray.toLong,dataArray.toInt,
dataArray,dataArray.toLong)
})
。分配取消时间戳
。过滤器
。地图)
。基比

。时间窗口)
。总额
。打印
环境执行
}
}
网站独立访问者数量的统计数据
在上面的例子中,我们统计了页面上所有用户的所有浏览行为,也就是说,同一个用户的浏览行为会被重复统计。在实践中,我们经常会关注一段时间内有多少不同的用户访问网站。
流量统计的另一个重要指标是网站的独立访问者数量。UV是指一段时间内访问网站的总人数。同一个访问者在一天内多次访问,只记录为一个访问者。一般IP和cookie是判断UV值的两种方式。当客户端第一次访问网站服务器时,网站服务器会向客户端计算机发送cookie,cookie通常放在客户端计算机的c盘中。在这个cookie中,会分配一个唯一的编号,这个编号会记录一些访问服务器的信息,比如访问时间,访问了哪些页面等等。当你下次访问这个服务器时,服务器可以直接从你的电脑中找到最后一个cookie文件并进行更新,但唯一编号不会改变。
当然,对于UserBehavior数据源,我们可以直接根据userId来区分不同的用户。
在src/main/scala下创建UniqueVisitor.scala文件,具体代码如下:
networkflow analysis/src/main/Scala/unique visitor . Scala
案例类UvCount
对象唯一访问者{
def main: Unit = {
val resources path = getclass . get resource
val env = streamexecutionenvironment . getexecutionenvironment
env . setstreamtimecharacter istic
env.setParallelism
val流=环境
。readTextFile
。地图
用户行为. toLong,linearray.toLong,linearray.toInt,
linearray,linearray.toLong)
})
。分配取消时间戳
。过滤器
。timeWindowAll)
。应用)
。打印
环境执行
}
}
UvCountByWindow类扩展了AllWindowFunction[UserBehavior,UvCount,TimeWindow] {
覆盖定义应用:单位= {
val s:collection . mutable . set[Long]= collection . mutable . set
var idSet = Set[Long]
对于{
idSet += userBehavior.userId
}
out.collect)
}
}
使用布隆过滤器的UV统计
在上一节的例子中,我们将所有数据的userId置于窗口计算的状态,在窗口数据采集的过程中,该状态会不断增加。一般情况下,只要不超出内存的范围,这个是没问题的。但是如果我们遇到大量的数据呢?
将所有数据临时存储在内存中显然不是一个好主意。我们会认为可以利用redis,一个内存级别的k-v数据库,为我们做一个缓存。但是如果我们遇到的情况非常极端,数据惊人呢?比如上亿用户要重新计算UV。
如果放在redis中,一个十亿级的用户id可能需要几个甚至几十个G 空来存储。当然,在redis中,用集群扩展也不是不可以,但显然太贵了。
更好的想法是,其实我们不需要完全存储用户ID的信息,只要知道他在不在就可以了。其实我们可以压缩一下,一个比特可以代表一个用户的状态。这一思想的具体实现就是Bloom filter。
本质上,Bloom filter是一种数据结构,一种巧妙的概率数据结构,特点是高效的插入和查询,可以用来告诉你“某个东西一定不存在或者可能存在”。
这是一个很长的二元向量。既然是二进制向量,显然存储的不是0就是1。与传统的数据结构如链表、集合、映射等相比。,它的效率更高,在空之间占用的空间更少,但它的缺点是返回的结果是概率性的,而不是精确的。
我们的目标是使用某种方法将每个数据对应到位图的一位;如果数据存在,该位为1,如果数据不存在,则为0。
接下来,我们来具体实施一下。
这里注意,我们使用redis连接来访问数据,所以我们需要加入redis客户端的依赖关系:
redis .客户
使用
2.8.1
在src/main/scala下创建UniqueVisitor.scala文件,具体代码如下:
networkflow analysis/src/main/Scala/uvwithloom . Scala
对象UvWithBloomFilter {
def main: Unit = {
val env = streamexecutionenvironment . getexecutionenvironment
env . setstreamtimecharacter istic
env.setParallelism
val resources path = getclass . get resource
val流=环境
。readTextFile
。地图
UserBehavior.toLong,dataArray.toLong,dataArray.toInt,
dataArray,dataArray.toLong)
})
。分配取消时间戳
。过滤器
。地图)
。基比
。时间窗口)
。trigger) //自定义窗口触发规则
。process) //自定义窗口处理规则
流.打印
环境执行
}
}
//自定义触发器
MyTrigger类扩展触发器[,TimeWindow] {
覆盖定义事件时间:
TriggerResult = {
触发结果。继续
}
覆盖def on processing time:trigger result = {
触发结果。继续
}
覆盖定义清除:单位= {
}
覆盖定义一个元素,时间戳:长整型,窗口:时间窗口,
ctx:触发。trigger context):trigger result = {
//每来一条数据,就触发窗口操作,清零空
触发结果。点火和吹扫
}
}
//自定义窗口处理函数
UvCountWithBloom类扩展了ProcessWindowFunction[,UvCount,String,
时间窗口] {
//创建redis连接
lazy val jedis =新jedis
懒惰的瓦尔布卢姆=新的布卢姆
覆盖定义流程],
out: Collector[UvCount]): Unit = {
val store key = context . window . getend . tostring
var count = 0L
如果!= null) {
count = jedis.hget.toLong
}
val userId = elements . last . _ 2 . tostring
val offset = bloom.hash
val isExist = jedis.getbit
如果{
jedis.setbit
jedis.hset.toString)
out.collect)
}否则{
out.collect)
}
}
}
//定义一个布隆过滤器
类Bloom扩展Serializable {
私有值上限=大小
定义哈希:Long = {
var结果= 0
对于{
//最简单的哈希算法,每个字符的ascii码值,乘以seed,叠加。

结果=结果*种子+值. charAt
}
结果
}
}


