最新国产好看的视频,伊人天堂AV在线,国产Aaaaaa视频,蜜臀视频在线观看一区,人妻av色图,密臀久久久精品影片,青青视频免费观看毛片,久草在线观看视,国产三级精品色情在线

Spark網站日志過濾分析實例講解

 更新時間:2023年02月01日 11:27:55   作者:CarveStone  
這篇文章主要介紹了Spark網站日志過濾分析實例,文中通過示例代碼介紹的非常詳細,對大家的學習或者工作具有一定的參考學習價值,需要的朋友們下面隨著小編來一起學習吧

日志過濾

對于一個網站日志,首先要對它進行過濾,刪除一些不必要的信息,我們通過scala語言來實現(xiàn),清洗代碼如下,代碼要通過別的軟件打包為jar包,此次實驗所用需要用到的代碼都被打好jar包,放到了/root/jar-files文件夾下:

package com.imooc.log
import com.imooc.log.SparkStatFormatJob.SetLogger
import com.imooc.log.util.AccessConvertUtil
import org.apache.spark.sql.{SaveMode, SparkSession}
/*
數(shù)據(jù)清洗部分
*/
object SparkStatCleanJob {
  def main(args: Array[String]): Unit = {
    SetLogger
    val spark = SparkSession.builder()
      .master("local[2]")
      .appName("SparkStatCleanJob").getOrCreate()
    val accessRDD = spark.sparkContext.textFile("/root/resources/access.log")
    accessRDD.take(4).foreach(println)
    val accessDF = spark.createDataFrame(accessRDD.map(x => AccessConvertUtil.parseLog(x)),AccessConvertUtil.struct)
    accessDF.printSchema()
    //-----------------數(shù)據(jù)清洗存儲到目標地址------------------------
    // coalesce(1)輸出指定分區(qū)數(shù)的小文件
    accessDF.coalesce(1).write.format("parquet").partitionBy("day").mode(SaveMode.Overwrite).save("/root/clean")//mode(SaveMode.Overwrite)覆蓋已經存在的文件  存儲為parquet格式,按day分區(qū)
      //存儲為parquet格式,按day分區(qū)
    /**
      * 調優(yōu)點:
      *     1) 控制文件輸出的大?。?coalesce
      *     2) 分區(qū)字段的數(shù)據(jù)類型調整:spark.sql.sources.partitionColumnTypeInference.enabled
      *     3) 批量插入數(shù)據(jù)庫數(shù)據(jù),提交使用batch操作
      */
    spark.stop()
  }
}

過濾好的數(shù)據(jù)將被存放在/root/clean文件夾中,這部分已被執(zhí)行好,后面直接使用就可以,其中代碼開始的SetLogger功能在自定義類com.imooc.log.SparkStatFormatJob中,它關閉了大部分log日志輸出,這樣可以使界面變得簡潔,代碼如下:

def SetLogger() = {
    Logger.getLogger("org").setLevel(Level.OFF)
    Logger.getLogger("com").setLevel(Level.OFF)
    System.setProperty("spark.ui.showConsoleProgress", "false")
    Logger.getRootLogger().setLevel(Level.OFF);
  }

過濾中的AccessConvertUtil類內容如下所示:

object AccessConvertUtil {
  //定義的輸出字段
  val struct = StructType(            //過濾日志結構
    Array(
      StructField("url", StringType), //課程URL
      StructField("cmsType", StringType), //課程類型:video / article
      StructField("cmsId", LongType), //課程編號
      StructField("traffic", LongType), //耗費流量
      StructField("ip", StringType), //ip信息
      StructField("city", StringType), //所在城市
      StructField("time", StringType), //訪問時間
      StructField("day", StringType) //分區(qū)字段,天
    )
  )
  /**
    * 根據(jù)輸入的每一行信息轉換成輸出的樣式
    * 日志樣例:2017-05-11 14:09:14     http://www.imooc.com/video/4500    304    218.75.35.226
    */
  def parseLog(log: String) = {
    try {
      val splits = log.split("\t")
      val url = splits(1)
      //http://www.imooc.com/video/4500
      val traffic = splits(2).toLong
      val ip = splits(3)
      val domain = "http://www.imooc.com/"
      //主域名
      val cms = url.substring(url.indexOf(domain) + domain.length)    //建立一個url的子字符串,它將從domain出現(xiàn)時的位置加domain的長度的位置開始計起
      val cmsTypeId = cms.split("/") 
      var cmsType = ""
      var cmsId = 0L
      if (cmsTypeId.length > 1) {
        cmsType = cmsTypeId(0)
        cmsId = cmsTypeId(1).toLong
      }      //以"/"分隔開后,就相當于分開了課程格式和id,以http://www.imooc.com/video/4500為例,此時cmsType=video,cmsId=4500
      val city = IpUtils.getCity(ip)         //從ip表中可以知道ip對應哪個城市
      val time = splits(0)
      //2017-05-11 14:09:14
      val day = time.split(" ")(0).replace("-", "")    //day=20170511
      //Row中的字段要和Struct中的字段對應
      Row(url, cmsType, cmsId, traffic, ip, city, time, day)
    } catch {
      case e: Exception => Row(0)
    }
  }
  def main(args: Array[String]): Unit = {
      //示例程序:
    val url = "http://www.imooc.com/video/4500"
    val domain = "http://www.imooc.com/" //主域名
    val index_0 = url.indexOf(domain)
    val index_1 = index_0 + domain.length
    val cms = url.substring(index_1)
    val cmsTypeId = cms.split("/")
    var cmsType = ""
    var cmsId = 0L
    if (cmsTypeId.length > 1) {
      cmsType = cmsTypeId(0)
      cmsId = cmsTypeId(1).toLong
    }
    println(cmsType + "   " + cmsId)
    val time = "2017-05-11 14:09:14"
    val day = time.split(" ")(0).replace("-", "")
    println(day)
  }
}

執(zhí)行完畢后clean文件夾下內容如圖1所示:

日志分析

現(xiàn)在我們已經擁有了過濾好的日志文件,可以開始編寫分析代碼,例如實現(xiàn)一個按地市統(tǒng)計主站最受歡迎的TopN課程

package com.imooc.log
import com.imooc.log.SparkStatFormatJob.SetLogger
import com.imooc.log.dao.StatDAO
import com.imooc.log.entity.{DayCityVideoAccessStat, DayVideoAccessStat, DayVideoTrafficsStat}
import org.apache.spark.sql.expressions.Window
import org.apache.spark.sql.functions._
import org.apache.spark.sql.{DataFrame, SparkSession}
import scala.collection.mutable.ListBuffer
object TopNStatJob2 {
  def main(args: Array[String]): Unit = {
    SetLogger
    val spark = SparkSession.builder()
      .config("spark.sql.sources.partitionColumnTypeInference.enabled", "false") //分區(qū)字段的數(shù)據(jù)類型調整【禁用】
      .master("local[2]")
      .config("spark.sql.parquet.compression.codec","gzip")   //修改parquet壓縮格式
      .appName("SparkStatCleanJob").getOrCreate()
    //讀取清洗過后的數(shù)據(jù)
    val cleanDF = spark.read.format("parquet").load("/root/clean")
    //執(zhí)行業(yè)務前先清空當天表中的數(shù)據(jù)
    val day = "20170511"
    import spark.implicits._
    val commonDF = cleanDF.filter($"day" === day && $"cmsType" === "video")
    commonDF.cache()
    StatDAO.delete(day)
    cityAccessTopSata(spark, commonDF)     //按地市統(tǒng)計主站最受歡迎的TopN課程功能
    commonDF.unpersist(true)     //RDD去持久化,優(yōu)化內存空間
    spark.stop()
  }
/*
 * 按地市統(tǒng)計主站最受歡迎的TopN課程
*/
def cityAccessTopSata(spark: SparkSession, commonDF: DataFrame): Unit = {
    //------------------使用DataFrame API完成統(tǒng)計操作--------------------------------------------
    import spark.implicits._
    val cityAccessTopNDF = commonDF
      .groupBy("day", "city", "cmsId").agg(count("cmsId").as("times")).orderBy($"times".desc)     //聚合
        cityAccessTopNDF.printSchema()
        cityAccessTopNDF.show(false)
     //-----------Window函數(shù)在Spark SQL中的使用--------------------
    val cityTop3DF = cityAccessTopNDF.select(       //Top3中涉及到的列
      cityAccessTopNDF("day"),
      cityAccessTopNDF("city"),
      cityAccessTopNDF("cmsId"),
      cityAccessTopNDF("times"),
      row_number().over(Window.partitionBy(cityAccessTopNDF("city"))
        .orderBy(cityAccessTopNDF("times").desc)).as("times_rank")
    ).filter("times_rank <= 3").orderBy($"city".desc, $"times_rank".asc)         //以city為一個partition,聚合times為times_rank,過濾出前三,降序聚合city,升序聚合times_rank
    cityTop3DF.show(false) //展示每個地市的Top3
     //-------------------將統(tǒng)計結果寫入數(shù)據(jù)庫-------------------
    try {
      cityTop3DF.foreachPartition(partitionOfRecords => {
        val list = new ListBuffer[DayCityVideoAccessStat]
        partitionOfRecords.foreach(info => {        
          val day = info.getAs[String]("day")
          val cmsId = info.getAs[Long]("cmsId")
          val city = info.getAs[String]("city")
          val times = info.getAs[Long]("times")
          val timesRank = info.getAs[Int]("times_rank")
          list.append(DayCityVideoAccessStat(day, cmsId, city, times, timesRank))
        })
        StatDAO.insertDayCityVideoAccessTopN(list)
      })
    } catch {
      case e: Exception => e.printStackTrace()
    }
    }

其中保存統(tǒng)計時用到了StatDAO類的insertDayCityVideoAccessTopN()方法,這部分的說明如下:

def insertDayVideoTrafficsTopN(list: ListBuffer[DayVideoTrafficsStat]): Unit = {
    var connection: Connection = null
    var pstmt: PreparedStatement = null
    try {
      connection = MySQLUtils.getConnection()      //JDBC連接MySQL
      connection.setAutoCommit(false) //設置手動提交
        //向day_video_traffics_topn_stat表中插入數(shù)據(jù)
      val sql = "insert into day_video_traffics_topn_stat(day,cms_id,traffics) values(?,?,?)"         
      pstmt = connection.prepareStatement(sql)
      for (ele <- list) {
        pstmt.setString(1, ele.day)
        pstmt.setLong(2, ele.cmsId)
        pstmt.setLong(3, ele.traffics)
        pstmt.addBatch() //優(yōu)化點:批量插入數(shù)據(jù)庫數(shù)據(jù),提交使用batch操作
      }
      pstmt.executeBatch() //執(zhí)行批量處理
      connection.commit() //手工提交
    } catch {
      case e: Exception => e.printStackTrace()
    } finally {
      MySQLUtils.release(connection, pstmt)          //釋放連接
    }
  }

JDBC連接MySQL和釋放連接用到了MySQLUtils中的方法

此外我們還需要在MySQL中插入表,用來寫入統(tǒng)計數(shù)據(jù),MySQL表已經設置好。

下面將程序和所有依賴打包,用spark-submit提交:

./spark-submit --class com.imooc.log.TopNStatJob2 --master spark://localhost:9000 /root/jar-files/sql-1.0-jar-with-dependencies.jar

執(zhí)行結果:

Schema信息

TopN課程信息

各地區(qū)Top3課程信息

MySQL表中數(shù)據(jù):

到此這篇關于Spark網站日志過濾分析實例講解的文章就介紹到這了,更多相關Spark日志分析內容請搜索腳本之家以前的文章或繼續(xù)瀏覽下面的相關文章希望大家以后多多支持腳本之家!

相關文章

  • Springboot整合EasyExcel實現(xiàn)Excel文件上傳方式

    Springboot整合EasyExcel實現(xiàn)Excel文件上傳方式

    這篇文章主要介紹了Springboot整合EasyExcel實現(xiàn)Excel文件上傳方式,具有很好的參考價值,希望對大家有所幫助,如有錯誤或未考慮完全的地方,望不吝賜教
    2024-07-07
  • Java生成隨機數(shù)的2種示例方法代碼

    Java生成隨機數(shù)的2種示例方法代碼

    在Java中,生成隨機數(shù)有兩種方法。1是使用Random類。2是使用Math類中的random方法??聪旅娴睦邮褂冒?/div> 2013-11-11
  • SpringMVC常用注解載入與處理方式詳解

    SpringMVC常用注解載入與處理方式詳解

    這篇文章主要介紹了SpringMVC常用注解載入的方式和處理的方式,文中示例代碼介紹的非常詳細,具有一定的參考價值,感興趣的小伙伴們可以參考一下
    2022-09-09
  • 淺談Spring Cloud Netflix-Ribbon灰度方案之Zuul網關灰度

    淺談Spring Cloud Netflix-Ribbon灰度方案之Zuul網關灰度

    這篇文章主要介紹了淺談Spring Cloud Netflix-Ribbon灰度方案之Zuul網關灰度,想了解Ribbon灰度的同學可以參考下
    2021-04-04
  • 一篇文章弄懂Java8中的時間處理

    一篇文章弄懂Java8中的時間處理

    Java8以前Java處理日期、日歷和時間的方式一直為社區(qū)所詬病,將 java.util.Date設定為可變類型,以及SimpleDateFormat的非線程安全使其應用非常受限,下面這篇文章主要給大家介紹了關于Java8中時間處理的相關資料,需要的朋友可以參考下
    2022-01-01
  • Java的CGLIB動態(tài)代理深入解析

    Java的CGLIB動態(tài)代理深入解析

    這篇文章主要介紹了Java的CGLIB動態(tài)代理深入解析,CGLIB是強大的、高性能的代碼生成庫,被廣泛應用于AOP框架,它底層使用ASM來操作字節(jié)碼生成新的類,為對象引入間接級別,以控制對象的訪問,需要的朋友可以參考下
    2023-11-11
  • Java設計模式之外觀模式示例詳解

    Java設計模式之外觀模式示例詳解

    外觀模式為多個復雜的子系統(tǒng),提供了一個一致的界面,使得調用端只和這個接口發(fā)生調用,而無須關系這個子系統(tǒng)內部的細節(jié)。本文將通過示例詳細為大家講解一下外觀模式,需要的可以參考一下
    2022-08-08
  • Java線程同步實例分析

    Java線程同步實例分析

    這篇文章主要介紹了Java線程同步用法,實例分析了java中線程同步的相關實現(xiàn)技巧與注意事項,具有一定參考借鑒價值,需要的朋友可以參考下
    2015-07-07
  • Java遞歸算法遍歷部門代碼示例

    Java遞歸算法遍歷部門代碼示例

    這篇文章主要介紹了Java遞歸算法遍歷部門代碼示例,具有一定借鑒價值,需要的朋友可以參考下。
    2017-12-12
  • JVM中的守護線程示例詳解

    JVM中的守護線程示例詳解

    這篇文章主要給大家介紹了關于JVM中守護線程的相關資料,文中通過示例代碼介紹的非常詳細,對大家的學習或者工作具有一定的參考學習價值,需要的朋友們下面隨著小編來一起學習學習吧
    2019-01-01

最新評論

博客| 镇赉县| 南宫市| 上饶县| 九龙县| 宿迁市| 普定县| 左贡县| 茶陵县| 仁寿县| 依兰县| 南部县| 高平市| 晋江市| 黄浦区| 新乡市| 苍山县| 福海县| 曲水县| 漳州市| 茂名市| 南丹县| 思茅市| 蒲城县| 镇江市| 沙田区| 富平县| 邵阳县| 吴忠市| 公安县| 南康市| 桑日县| 商丘市| 静乐县| 大关县| 安塞县| 建平县| 吉林省| 邵东县| 大庆市| 来安县|