1.主题
使用DataSet API统计“1小时内”用户的访问链接次数(注意理解1小时内的含义)
2.适用读者对象
本文适用于hadoop大数据Flink入门学员
3.基础环境信息
| 软件 | 版本 | 部署路径 |
haoop | 2.6.0-cdh5.14.4 | cdh默认安装 |
flink | flink-1.11.1 | /home/opt/flink/flink-1.11.1 |
4.主题描述
DataSet API主要应用于有界数据集场景,用户可自定义转换(transformations)、联接(joins)、聚合(aggregations)、操作等。
本文使用DataSet API统计“1小时内”用户的访问链接次数,比如当统计数据(Liz,2018-12-22 14:02:00,./home)的时候,你需要统计用户名为Liz,时间在“2018-12-22 13:02:01-2018-12-22 14:02:00”的所有数据。由于是有界数据集,这里的输入和输出都是以批量形式处理的。
5.实战演练
5.1 数据处理流图

5.2 输入数据
user,time,url
Mary,2018-12-22 12:00:00,./home Bob,2018-12-22 12:00:00,./cart Mary,2018-12-22 12:02:00,./prod?id=1 Mary,2018-12-22 12:55:00,./prod?id=4 Bob,2018-12-22 13:01:00,./prod?id=5 Liz,2018-12-22 13:30:00,./home Liz,2018-12-22 13:59:00,./prod?id=7 Mary,2018-12-22 14:00:00,./cart Liz,2018-12-22 14:02:00,./home Bob,2018-12-22 14:30:00,./prod?id=3 Bob,2018-12-22 14:40:00,./home Bob,2018-12-22 15:40:00,./home |
5.3 代码实现
package com.zhenglihan.cdh.flinkimport java.text.SimpleDateFormatimport java.util.{Calendar}import org.apache.flink.api.common.functions.GroupReduceFunctionimport org.apache.flink.api.scala._import org.apache.flink.util.Collectorimport scala.collection.JavaConversions/*** 功能描述:* 统计“一小时内”用户的访问链接次数(注意理解一小时内的含义)* 实现方案:* 采用flink Dataset API 加自定义函数* 关键知识点:* GroupReduceFunction* 输入:* DataSet* 输出:* 控制台* 输出结果:* (Bob,2018-12-22 12:00:00,1)* (Bob,2018-12-22 13:01:00,1)* (Bob,2018-12-22 14:30:00,1)* (Bob,2018-12-22 14:40:00,2)* (Bob,2018-12-22 15:40:00,1)* (Liz,2018-12-22 13:30:00,1)* (Liz,2018-12-22 13:59:00,2)* (Liz,2018-12-22 14:02:00,3)* (Mary,2018-12-22 12:00:00,1)* (Mary,2018-12-22 12:02:00,2)* (Mary,2018-12-22 12:55:00,3)* (Mary,2018-12-22 14:00:00,1)*/object FlinkBatchDataSetUserClickDemo {def main(args: Array[String]): Unit = {// set up execution environmentval env = ExecutionEnvironment.getExecutionEnvironmentval input = env.fromElements(UserClick("Mary","2018-12-22 12:00:00","./home"),UserClick("Bob","2018-12-22 12:00:00","./cart"),UserClick("Mary","2018-12-22 12:02:00","./prod?id=1"),UserClick("Mary","2018-12-22 12:55:00","./prod?id=4"),UserClick("Bob","2018-12-22 13:01:00","./prod?id=5"),UserClick("Liz","2018-12-22 13:30:00","./home"),UserClick("Liz","2018-12-22 13:59:00","./prod?id=7"),UserClick("Mary","2018-12-22 14:00:00","./cart"),UserClick("Liz","2018-12-22 14:02:00","./home"),UserClick("Bob","2018-12-22 14:30:00","./prod?id=3"),UserClick("Bob","2018-12-22 14:40:00","./home"),UserClick("Bob","2018-12-22 15:40:00","./home"))input.groupBy(_.user).reduceGroup(new UserClickGroupReduceFunction).print()}case class UserClick(user: String, time: String, url:String)class UserClickGroupReduceFunction extends GroupReduceFunction[UserClick, (String,String,Long)] {override def reduce(input: java.lang.Iterable[UserClick], out: Collector[(String,String,Long)]) = {val list=JavaConversions.iterableAsScalaIterable(input).toList.sortBy(_.time)list.foreach(userClick=>{var count=0for(i<-0 to list.size-1){val userClick2=list(i)//统计list中userClick.time数据时间在一小时以内的数据if(userClick2.time.compareTo(getOneHourBefore(userClick.time,-1)+userClick.time.substring(13))>0 && userClick2.time.compareTo(userClick.time)<=0){count=count+1}}out.collect((userClick.user,userClick.time,count))})out}/*** 获取preHour小时后的时间,精确到小时,例如输入:2018-12-22 13:59:00,-1 输出:2018-12-22 12* @param time* @param preHour* @return*/def getOneHourBefore(time:String,preHour:Int): String = {val cal = Calendar.getInstanceval sdf = new SimpleDateFormat("yyyy-MM-dd HH");cal.setTime(sdf.parse(time))cal.add(Calendar.HOUR_OF_DAY, preHour)val newTime = sdf.format(cal.getTime)newTime}}}
5.4 运行结果

6.总结:
本文使用DataSet API统计“1小时内”用户的访问链接次数。以批处理的形式模拟实现流式处理的逻辑,核心代码在于reduceGroup和GroupReduceFunction的运行机制。
更为深度的讲解请扫描底部二维码关注公众号,关注后续博文,一起学习hadoop大数据!如果读者有啥疑问或者建议,也可以在底部评论,留言,建议,我会尽力帮大家解决疑惑!






