spark企业级电商分析平台项目实践(四)需求一分析与实现

    科技2026-10-02  8

    前言

    需求一:各范围session步长、时长所占比例统计的分析与实现

    一、需求

    在上一章节,根据过滤条件,已经得到了(sessionId, filteredRDD),并且更新了累加器中的三个field:1.session_count 2.time_period 3.step_period,根据累加器中数据,计算session各步长、各时长所占比例就可以了。

    二、实现

    声明 //需求一:各范围session占比统计 getSessionRatio(sparkSession, taskUUID, sessionStatisticAccumulator.value) 定义 def getSessionRatio(sparkSession: SparkSession, taskUUID: String, value: mutable.HashMap[String, Int]) :Unit={ //获取sessionCount作为分母,防止分母为零,给个默认1 val session_count = value.getOrElse(Constants.SESSION_COUNT, 1).toDouble //获取各步长、时长的值作为分子,hashmap用getOrElse取出值 val visit_length_1s_3s = value.getOrElse(Constants.TIME_PERIOD_1s_3s, 0) ... val step_length_1_3 = value.getOrElse(Constants.STEP_PERIOD_1_3, 0) ... //取值完毕后,与sessionCount相除,按指定格式保留小数位数 val visit_length_1s_3s_ratio = NumberUtils.formatDouble(visit_length_1s_3s / session_count, 2) ... val step_length_1_3_ratio = NumberUtils.formatDouble(step_length_1_3 / session_count, 2) ... //获得所有占比后,将各个占比封装case class val stat = SessionAggrStat(taskUUID, ... ,step_length_60_ratio) // 把对象封装成一个rdd,再转成df,写入mysql val statRDD = sparkSession.sparkContext.makeRDD(Array(stat)) import sparkSession.implicits._ statRDD.toDF().write .format("jdbc") .option("url", ConfigurationManager.config.getString(Constants.JDBC_URL)) .option("user", ConfigurationManager.config.getString(Constants.JDBC_USER)) .option("password", ConfigurationManager.config.getString(Constants.JDBC_PASSWORD)) .option("dbtable", "session_stat_1007") .mode(SaveMode.Append) //保存模式为追加,否则每条数据创建表,会显示表已存在 .save() }
    Processed: 0.009, SQL: 10