暂无图片
暂无图片
暂无图片
暂无图片
暂无图片

hadoop大数据Flink专题篇-Flink四种API入门实战(二)

hadoop大数据入门引路人 2020-09-01
421

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.flink


    import java.text.SimpleDateFormat
    import java.util.{Calendar}


    import org.apache.flink.api.common.functions.GroupReduceFunction
    import org.apache.flink.api.scala._
    import org.apache.flink.util.Collector


    import 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 environment
    val env = ExecutionEnvironment.getExecutionEnvironment


    val 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=0
    for(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.getInstance
    val 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小时内”用户的访问链接次数。以批处理的形式模拟实现流式处理的逻辑,核心代码在于reduceGroupGroupReduceFunction的运行机制

        

        更为深度的讲解请扫描底部二维码关注公众号,关注后续博文一起学习hadoop大数据!如果读者有啥疑问或者建议,也可以在底部评论,留言,建议,我会尽力帮大家解决疑惑!

    文章转载自hadoop大数据入门引路人,如果涉嫌侵权,请发送邮件至:contact@modb.pro进行举报,并提供相关证据,一经查实,墨天轮将立刻删除相关内容。

    评论