idmapping(用户唯一标识)真实数据第二天数据生成-程序员宅基地

技术标签: 大数据  

idmapping(用户唯一标识)第二天数据生成

/**
  *思想逻辑:(必须整明白!!!)
  * 我们是要考虑第二天的数据进来了
  * 在拿去跟第一天的数据一起进行计算
  * 我们需要改的地方就是在构造图的地方添加昨天的点集合和边集合
  *
  * 1.将昨天的字典解析成点集合,边集合,
  * 2.将今天的点集合,边集合union到昨天的点边集合
  * 3.用union之后的点边集合构造最大连同子图
  * 4.调整结果
  *   今天有一套结果,昨天也有一套结果,
  *   虽然两套结果都对,但是总有可能今天的结果新生成的guid和昨天的有一点点小误差,
  *   为了方便日后留存查询,我们将今天有小误差的guid改成昨天的
  *   只有这样才能保证今天生成的guid和昨天生的guid以及以后所有的guid才能使同一个人
  *   这样我们以后查询当月的人流量才能确定是同一个人,避免泡沫数据和虚假数据!!!!
  *   步骤实现:
  *   4.1:我们将今天guid的行组合在一起,将昨天的字典搞成一个hashmap进行广播
  *   4.1:拿今天字典的点集合和昨天的额点集合求交集取到谁今天的guid就换成昨天的guid
  *
  */
object idmapping_nextday {
    
  def main(args: Array[String]): Unit = {
    
    //1.1.创建spark环境
    val spark: SparkSession = SparkSession.builder()
      .appName(this.getClass.getSimpleName)
      .master("local[*]")
      .getOrCreate()
    //将rdd变成df
    import spark.implicits._
    //2.导入数据
    val applog: Dataset[String] = spark.read.textFile("F:\\yiee_logs\\2020-01-11\\app\\doit.mall.access.log.8")

    //3.处理数据
    //构造一个点集合
    val data: RDD[Array[String]] = applog.rdd.map(line => {
    
      //将每行数据解析成json对象
      val jsonObj = JSON.parseObject(line)

      // 从json对象中取user对象
      val userObj = jsonObj.getJSONObject("user")
      val uid = userObj.getString("uid")

      // 从user对象中取phone对象
      val phoneObj = userObj.getJSONObject("phone")
      val imei = phoneObj.getString("imei")
      val mac = phoneObj.getString("mac")
      val imsi = phoneObj.getString("imsi")
      val androidId = phoneObj.getString("androidId")
      val deviceId = phoneObj.getString("deviceId")
      val uuid = phoneObj.getString("uuid")

      Array(uid, imei, mac, imsi, androidId, deviceId, uuid).filter(StringUtils.isNotBlank(_))
    })

    //构造一个点集合
    // 将每一个元素和他的hashcode组合成对偶元组
    val vertices: RDD[(Long, String)] = data.flatMap(arr => {
    
      for (biaoshi <- arr) yield (biaoshi.hashCode.toLong, biaoshi)
    })

    //构造图计算的边集合
    // Edge参数一:第一个点.参数二:下一个点.参数三:边的数据或者边的属性,这里我们用空串表示
    //将两个点连接起来
    val edges: RDD[Edge[String]] = data.flatMap(arr => {
    
      //用双重for循环的方法让数组中所有的两两组合成边
      for (i <- 0 to arr.length - 2; j <- i + 1 to arr.length - 1) yield Edge(arr(i).hashCode.toLong, arr(j).hashCode.toLong, "")
    })
      //然后在统计每一条边出现的次数
      .map(edge => (edge, 1)).reduceByKey(_ + _)
      //过滤将重复次数<5(经验阈值)的边去掉,
      .filter(tp => tp._2 > 2)
      .map(x => x._1)
    edges



    // 五、将上一日的idmp映射字典,解析成点、边集合
    val schema = new StructType()
      .add("biaoshi_hashcode",DataTypes.LongType)
      .add("guid",DataTypes.LongType)


    //新0将昨天的节过读出来
    val preDayIdmp: DataFrame = spark.read.parquet("data/dict/idmapping_dict/2020-02-15")
    //新1.读取昨天的字典解析成点集合
    val preDayIdmpVertices = preDayIdmp.rdd.map({
    
      case Row(idFlag: VertexId, guid: VertexId) =>
        (idFlag, "")
    })

    //新2.读取昨天的字典解析成边集合
    val preDayEdges = preDayIdmp.rdd.map(row => {
    
      val idFlag = row.getAs[VertexId]("biaoshi_hashcode")
      val guid = row.getAs[VertexId]("guid")
      Edge(idFlag, guid, "")
    })

    //新3.用 点集合 和 边集合 构造一张图  使用Graph算法和昨天的结果union在一起
    val graph = Graph(vertices.union(preDayIdmpVertices),edges.union(preDayEdges))

    //并调用最大连同子图算法VertexRDD[VertexId] ==>rdd 里面装的元组(Long值,组中最小值)
    val res_tuples: VertexRDD[VertexId] = graph.connectedComponents().vertices
    //最终我们要输出的文件是parquet的文件






    // 最后的调整结果、将结果跟上日的映射字典做对比,调整guid
    // 1.将上日的idmp映射结果字典收集到driver端,并广播
    val idMap = preDayIdmp.rdd.map(row => {
    
      val idFlag = row.getAs[VertexId]("biaoshi_hashcode")
      val guid = row.getAs[VertexId]("guid")
      (idFlag, guid)
    }).collectAsMap()
    val bc = spark.sparkContext.broadcast(idMap)

    // 2.将今日的图计算结果按照guid分组,然后去跟上日的映射字典进行对比
    val todayIdmpResult: RDD[(VertexId, VertexId)] = res_tuples.map(tp => (tp._2, tp._1))
      .groupByKey()
      .mapPartitions(iter=>{
    
        // 从广播变量中取出上日的idmp映射字典
        val idmpMap = bc.value
        iter.map(tp => {
    
          // 当日的guid计算结果
          var todayGuid = tp._1
          // 这一组中的所有id标识
          val ids = tp._2

          // 遍历这一组id,挨个去上日的idmp映射字典中查找
          var find = false
          for (elem <- ids if !find) {
    
            val maybeGuid: Option[Long] = idmpMap.get(elem)
            // 如果这个id在昨天的映射字典中找到了,那么就用昨天的guid替换掉今天这一组的guid
            if (maybeGuid.isDefined) {
    
              todayGuid = maybeGuid.get
              find = true
            }
          }

          (todayGuid,ids)
        })
      })
      .flatMap(tp=>{
    
        val ids = tp._2
        val guid = tp._1
        for (elem <- ids) yield (elem,guid)
      })
    // 可以直接用图计算所产生的结果中的组最小值,作为这一组的guid(当然,也可以自己另外生成一个UUID来作为GUID)
    import spark.implicits._
    // 保存结果
    todayIdmpResult.coalesce(1).toDF("biaoshi_hashcode", "guid").write.parquet("data/dict/idmapping_dict/2020-02-16")



    spark.stop()
  }
}
版权声明:本文为博主原创文章,遵循 CC 4.0 BY-SA 版权协议,转载请附上原文出处链接和本声明。
本文链接:https://blog.csdn.net/weixin_45896475/article/details/104333044

智能推荐

oracle 12c 集群安装后的检查_12c查看crs状态-程序员宅基地

文章浏览阅读1.6k次。安装配置gi、安装数据库软件、dbca建库见下:http://blog.csdn.net/kadwf123/article/details/784299611、检查集群节点及状态:[root@rac2 ~]# olsnodes -srac1 Activerac2 Activerac3 Activerac4 Active[root@rac2 ~]_12c查看crs状态

解决jupyter notebook无法找到虚拟环境的问题_jupyter没有pytorch环境-程序员宅基地

文章浏览阅读1.3w次,点赞45次,收藏99次。我个人用的是anaconda3的一个python集成环境,自带jupyter notebook,但在我打开jupyter notebook界面后,却找不到对应的虚拟环境,原来是jupyter notebook只是通用于下载anaconda时自带的环境,其他环境要想使用必须手动下载一些库:1.首先进入到自己创建的虚拟环境(pytorch是虚拟环境的名字)activate pytorch2.在该环境下下载这个库conda install ipykernelconda install nb__jupyter没有pytorch环境

国内安装scoop的保姆教程_scoop-cn-程序员宅基地

文章浏览阅读5.2k次,点赞19次,收藏28次。选择scoop纯属意外,也是无奈,因为电脑用户被锁了管理员权限,所有exe安装程序都无法安装,只可以用绿色软件,最后被我发现scoop,省去了到处下载XXX绿色版的烦恼,当然scoop里需要管理员权限的软件也跟我无缘了(譬如everything)。推荐添加dorado这个bucket镜像,里面很多中文软件,但是部分国外的软件下载地址在github,可能无法下载。以上两个是官方bucket的国内镜像,所有软件建议优先从这里下载。上面可以看到很多bucket以及软件数。如果官网登陆不了可以试一下以下方式。_scoop-cn

Element ui colorpicker在Vue中的使用_vue el-color-picker-程序员宅基地

文章浏览阅读4.5k次,点赞2次,收藏3次。首先要有一个color-picker组件 <el-color-picker v-model="headcolor"></el-color-picker>在data里面data() { return {headcolor: ’ #278add ’ //这里可以选择一个默认的颜色} }然后在你想要改变颜色的地方用v-bind绑定就好了,例如:这里的:sty..._vue el-color-picker

迅为iTOP-4412精英版之烧写内核移植后的镜像_exynos 4412 刷机-程序员宅基地

文章浏览阅读640次。基于芯片日益增长的问题,所以内核开发者们引入了新的方法,就是在内核中只保留函数,而数据则不包含,由用户(应用程序员)自己把数据按照规定的格式编写,并放在约定的地方,为了不占用过多的内存,还要求数据以根精简的方式编写。boot启动时,传参给内核,告诉内核设备树文件和kernel的位置,内核启动时根据地址去找到设备树文件,再利用专用的编译器去反编译dtb文件,将dtb还原成数据结构,以供驱动的函数去调用。firmware是三星的一个固件的设备信息,因为找不到固件,所以内核启动不成功。_exynos 4412 刷机

Linux系统配置jdk_linux配置jdk-程序员宅基地

文章浏览阅读2w次,点赞24次,收藏42次。Linux系统配置jdkLinux学习教程,Linux入门教程(超详细)_linux配置jdk

随便推点

matlab(4):特殊符号的输入_matlab微米怎么输入-程序员宅基地

文章浏览阅读3.3k次,点赞5次,收藏19次。xlabel('\delta');ylabel('AUC');具体符号的对照表参照下图:_matlab微米怎么输入

C语言程序设计-文件(打开与关闭、顺序、二进制读写)-程序员宅基地

文章浏览阅读119次。顺序读写指的是按照文件中数据的顺序进行读取或写入。对于文本文件,可以使用fgets、fputs、fscanf、fprintf等函数进行顺序读写。在C语言中,对文件的操作通常涉及文件的打开、读写以及关闭。文件的打开使用fopen函数,而关闭则使用fclose函数。在C语言中,可以使用fread和fwrite函数进行二进制读写。‍ Biaoge 于2024-03-09 23:51发布 阅读量:7 ️文章类型:【 C语言程序设计 】在C语言中,用于打开文件的函数是____,用于关闭文件的函数是____。

Touchdesigner自学笔记之三_touchdesigner怎么让一个模型跟着鼠标移动-程序员宅基地

文章浏览阅读3.4k次,点赞2次,收藏13次。跟随鼠标移动的粒子以grid(SOP)为partical(SOP)的资源模板,调整后连接【Geo组合+point spirit(MAT)】,在连接【feedback组合】适当调整。影响粒子动态的节点【metaball(SOP)+force(SOP)】添加mouse in(CHOP)鼠标位置到metaball的坐标,实现鼠标影响。..._touchdesigner怎么让一个模型跟着鼠标移动

【附源码】基于java的校园停车场管理系统的设计与实现61m0e9计算机毕设SSM_基于java技术的停车场管理系统实现与设计-程序员宅基地

文章浏览阅读178次。项目运行环境配置:Jdk1.8 + Tomcat7.0 + Mysql + HBuilderX(Webstorm也行)+ Eclispe(IntelliJ IDEA,Eclispe,MyEclispe,Sts都支持)。项目技术:Springboot + mybatis + Maven +mysql5.7或8.0+html+css+js等等组成,B/S模式 + Maven管理等等。环境需要1.运行环境:最好是java jdk 1.8,我们在这个平台上运行的。其他版本理论上也可以。_基于java技术的停车场管理系统实现与设计

Android系统播放器MediaPlayer源码分析_android多媒体播放源码分析 时序图-程序员宅基地

文章浏览阅读3.5k次。前言对于MediaPlayer播放器的源码分析内容相对来说比较多,会从Java-&amp;amp;gt;Jni-&amp;amp;gt;C/C++慢慢分析,后面会慢慢更新。另外,博客只作为自己学习记录的一种方式,对于其他的不过多的评论。MediaPlayerDemopublic class MainActivity extends AppCompatActivity implements SurfaceHolder.Cal..._android多媒体播放源码分析 时序图

java 数据结构与算法 ——快速排序法-程序员宅基地

文章浏览阅读2.4k次,点赞41次,收藏13次。java 数据结构与算法 ——快速排序法_快速排序法