十年匠心定制 · 商业建站与技术教学双线并行 咨询热线:400-886-1026 service@lmnt.cn
ARTICLE DETAIL

资讯详情

深耕网站建设与运营推广的一线实战洞察。

Spark电商用户行为分析系统实战:从日志清洗到留存率计算

Spark电商用户行为分析系统实战:从日志清洗到留存率计算 简介基于Spark的电商用户行为分析系统完整算法源码包面向毕业设计、课程设计及Spark技术学习者可快速搭建完整的用户行为分析项目。系统核心覆盖用户Session会话分析、区域热门商品Top3统计、需求函数分析等典型场景全部以Scala实现体现从原始日志清洗、业务聚合到结果输出的完整流程。压缩包共273个文件包括58个Scala源码文件、208个XML配置资源以及properties属性文件、README说明文档等整体体积仅183KB。其中XML文件用于工程配置与作业描述Scala文件按业务模块拆分方便按需阅读与调试。资源包轻量且依赖简单适合本地快速运行与验证截至当前已有125人学习下载可供需要参考Spark算子实际用法、完成课程设计或毕业设计功能验证的开发者直接使用。1. 基于Spark的电商用户行为分析系统到底在解决什么问题拿到一个基于Spark的电商用户行为分析系统.zip大多数人第一反应是解压看代码但真正值得先想清楚的是这套系统在业务上到底要回答哪几个问题。用户行为分析不是把PV/UV算出来而是把一次点击、一次加购、一次支付按时间串成行为序列再算漏斗、复购和留存。Spark在这里的角色不是包治百病的万能引擎而是把分散在日志文件和业务库里的海量行为数据清洗成一张可供数据仓库反复查询的宽表。这套系统适合数据分析师、数仓工程师和刚接手Spark业务的开发下面就是我搭建这样一个系统时会完整走一遍的路径。2. 电商用户行为分析系统的数据接入和日志字段设计2.1 数据源先统一schema还是先接进来再说电商行为日志的主源头是前端埋点常见的存储方式是按天分目录的JSON文件。订单数据一般躺在MySQL里商品维表则在Hive中。先约定一个最小schema再写Spark读入否则到了清洗环节会反复返工。我一般会要求埋点必须输出userId、itemId、categoryId、behaviorType、eventTime、extra。userId可以是未登录时由前端根据浏览器指纹生成的临时ID但要保证同一个会话内不变behaviorType只允许pv/cart/fav/buy四种枚举eventTime是毫秒级Unix时间戳注意不要混用秒级。这里还有一个经常被忽略的点日志采集链路中前端SDK上报后要经过Nginx和日志采集组件才能落盘数据流上任何一环做标准化失败都会让字段类型漂移。所以设计schema时userId和itemId都定义为String而不是Long避免前导0被丢弃categoryId允许空值。extra字段是JSON字符串会在清洗阶段解析。订单表字段准备orderId、userId、itemId、amount、paidTime用于离线对账和交易漏斗。下面是标准行为日志字段表字段名类型约束说明userIdString非空用户唯一标识itemIdString非空商品IDcategoryIdString可空类目IDbehaviorTypeString枚举(pv/cart/fav/buy)行为类型eventTimeLong毫秒事件时间戳extraStringJSON页面来源、页面名称等2.2 用DataFrameReader读取多种来源的标准姿势读取JSON时不能偷懒让Spark自动推断schema因为日志量一旦增大自动推断会额外扫描一整遍文件且碰到类型不统一时会直接报错。这里给出明确schema的读法from pyspark.sql.types import StructField, StructType, StringType, LongType behavior_schema StructType([ StructField(userId, StringType(), True), StructField(itemId, StringType(), True), StructField(categoryId, StringType(), True), StructField(behaviorType, StringType(), True), StructField(eventTime, LongType(), True), StructField(extra, StringType(), True) ]) raw_df spark.read.schema(behavior_schema) \ .option(badRecordsPath, /tmp/behavior_bad_records) \ .json(/data/logs/events/dt2024-06-01)这样Spark不需要预扫描遇到无法解析的字段会按null处理同时把坏记录写到指定路径方便排查。接下来读MySQL订单表重点要设分区字段order_df spark.read \ .format(jdbc) \ .option(url, jdbc:mysql://10.0.2.10:3306/ecom?useSSLfalse) \ .option(dbtable, orders) \ .option(user, spark_reader) \ .option(password, ******) \ .option(partitionColumn, id) \ .option(lowerBound, 1) \ .option(upperBound, 10000000) \ .option(numPartitions, 8) \ .load()partitionColumn必须是数值列Spark用它做并行范围切分如果不设JDBC会全表扫描拉回Driver数据库压力很大。同时建议在MySQL侧为where条件建索引比如paid_time。注意不要把dbtable写成子查询比如(select * from orders where paid_time ...) as t这会让Spark把整表拉回内存再过滤失去下推效果。读Hive维表就更简单了直接用SparkSQLitem_df spark.sql(SELECT itemId, categoryId, title, price FROM dwd.dim_item WHERE dt 2024-06-01)这里的关键是必须在SQL里写分区条件让Spark走分区裁剪否则全表扫描一个维表会让后续join都变慢。2.3 为什么偏偏是DataFrame而不是RDD如果只会RDD也可以写行为分析但你要手动做分区器、广播变量和combiner。DataFrame自带Catalyst优化器它会对过滤条件下推对列裁剪甚至把多个相邻算子重排成更高效的执行计划。尤其是在Spark 2.0之后DataFrameDataset[Row]成了标准APISpark SQL能直接复用它。同一套代码既能跑在本地又能用Spark on YARN提交到集群。“Spark之DataFrame”在面试或实际代码评审中经常被问我一般会强调三点一、写业务逻辑比RDD短一半以上二、执行计划可以用explain查看瓶颈可见三、在Spark SQL和DataFrame API之间可以无缝切换。用户行为分析涉及大量groupBy、窗口和joinDataFrame天然适合比如执行计划里能看到Filter被下推到数据源这通常是优化必须验确认的一环。3. 用Spark SQL做行为清洗与会话切分的核心实现3.1 清洗不是简单的filter三个动作必须做行为数据清洗至少要完成重复上报处理、无主用户过滤、时间字段标准化。前端在弱网环境会重发事件同一用户在1秒内对同一商品做同一行为基本可以判定为重复。使用row_number加窗口去重from pyspark.sql.window import Window from pyspark.sql.functions import col, row_number, from_unixtime dedup_spec Window.partitionBy(userId, itemId, behaviorType, eventTime) \ .orderBy(eventTime) clean_df raw_df \ .withColumn(rn, row_number().over(dedup_spec)) \ .filter(rn 1) \ .drop(rn) \ .filter(userId IS NOT NULL AND eventTime IS NOT NULL) \ .withColumn(event_time, from_unixtime(col(eventTime) / 1000, yyyy-MM-dd HH:mm:ss))这样把重复事件去掉后真正时间字段也生成出来了。字段event_time后续会作为会话窗口排序和前一日分区关联的键。清洗阶段尽量避免用自定义UDF做正则过滤能用内置函数就用内置函数原因很简单内置函数走的是代码生成路径性能比Python UDF高一个数量级。另外要明确清洗范围如果本次分析只关注近30天的数据就在读取时一次性过滤比如raw_df.filter(eventTime start_ts AND eventTime end_ts)不要让后面每个计算步骤都背着无用的历史数据。我遇到过因为多跑了两个月日志导致shuffle数据量翻了5倍的例子。3.2 会话切分的30分钟阈值用lag和sum实现会话指用户在限定时间内连续访问商品页的过程。超过30分钟没有行为下一次行为就算新会话。阈值业界常用但严格说应该用业务方的平均停留时间标定。实现方法先按userId分区用eventTime排序再取出上一个事件时间差值大于阈值则打标记1最后对这个标记求和得到每个事件的sessionId。代码from pyspark.sql.functions import lag, when, sum session_window Window.partitionBy(userId).orderBy(eventTime) sessioned_df clean_df \ .withColumn(prev_time, lag(eventTime).over(session_window)) \ .withColumn(new_session, when(col(prev_time).isNull(), 1) .when((col(eventTime) - col(prev_time)) 1800000, 1) .otherwise(0)) \ .withColumn(session_gid, sum(new_session).over(session_window))这里的1800000毫秒就是30分钟。如果未来想对比不同阈值的效果把这个值抽成参数即可。注意一个坑如果原始eventTime存在乱序比如埋点上报延后切分会出错所以在进入切分前最好按userIdeventTime做一次全局排序或者让上游写日志时就按时间写入。会话切分后还能顺带算出会话时长和页面深度这些指标后面做用户路径分析时很管用。3.3 数据倾斜处理改一个分区数不够用户行为数据有一个天然特征少数头部用户的日志量远超普通用户。直接对userId做聚合会出现多个executor空转、一个task拖垮整个stage。常见处理有三个层次。第一先调整shuffle参数这是最省事的优化--conf spark.sql.shuffle.partitions240 --conf spark.sql.adaptive.enabledtrue --conf spark.sql.adaptive.skewJoin.enabledtrueSpark 3.0之后的adaptive执行可以在shuffle结束后自动合并小分区和拆分倾斜分区不一定要改代码。第二对确实热点明显的key加盐。下面是一个对用户行为数做加盐统计的示意from pyspark.sql.functions import col, rand, split, sum, count salted sessioned_df \ .withColumn(salt, (rand() * 10).cast(int)) \ .withColumn(salted_user, concat(col(userId), lit(_), col(salt))) \ .groupBy(salted_user).agg(count(*).alias(partial_cnt)) \ .withColumn(userId, split(salted_user, _)[0]) \ .groupBy(userId).agg(sum(partial_cnt).alias(behavior_cnt))先随机打散成10份局部聚合成partial_cnt再按真实userId求和能把热点用户分到多个task。加盐数不是越大越好10到20份通常够用。第三在构造宽表时把热点维表做成广播变量避免大join。比如商品维表只有几十万条就用from pyspark.sql.functions import broadcast result_df sessioned_df.join(broadcast(item_df), itemId, left)这样join会走BroadcastHashJoin不会触发全shuffle。关于“spark集群搭建”后的资源问题调参前先看executor数量和并行度是否匹配并行度一般设为总核心数的2到3倍而不是盲目调大。4. 用户行为分析核心指标漏斗、复购率与留存率4.1 漏斗转化用DataFrame的聚合加JOIN不用反复count漏斗分析要给出从浏览到支付的每一步用户数。直接从行为日志里分步count最怕一个用户跨多天被重复计数。正确做法是在明细数据里为每个userId生成一个是否发生某行为的标志列然后再聚合from pyspark.sql.functions import max, when funnel_flag sessioned_df.groupBy(userId).agg( max(when(col(behaviorType) pv, 1).otherwise(0)).alias(has_pv), max(when(col(behaviorType) cart, 1).otherwise(0)).alias(has_cart), max(when(col(behaviorType) fav, 1).otherwise(0)).alias(has_fav), max(when(col(behaviorType) buy, 1).otherwise(0)).alias(has_buy) ) funnel_step funnel_flag.selectExpr( COUNT(IF(has_pv1,1,NULL)) AS pv_users, COUNT(IF(has_cart1,1,NULL)) AS cart_users, COUNT(IF(has_buy1,1,NULL)) AS buy_users )用max聚合是因为一个人只要发生过加购就置1。计算步骤间比率时从pv到cart的转化率就是cart_users/pv_users。注意执行顺序如果先筛选只用buy用户再算加购转化就会丢掉中途流失的信息。如果想看每一步之间的流失还要在四个标志列之间做交叉统计这时候SQL的case when比DataFrame里多次join更直观。4.2 复购率计算先按用户和日期去重复购率通常定义为一个周期内下单超过一次的用户占比。由于同一用户当天可能下多笔单不能直接把订单数当购买天数否则会虚高。我一般会先把订单按userId和日期去重得到“用户每日购买事实”再按用户统计购买天数from pyspark.sql.functions import col, to_date buy_daily order_df \ .filter(paid_time IS NOT NULL) \ .select(col(userId), to_date(col(paid_time)).alias(buy_date)) \ .distinct() buy_stats buy_daily.groupBy(userId).count().toDF(userId, buy_days) reorder_rate buy_stats.selectExpr( SUM(IF(buy_days 2, 1, 0)) AS reorder_users, COUNT(*) AS active_users ).selectExpr(reorder_users / active_users AS reorder_rate)计算周期参数也要想清楚统计2024年6月的复购率就要把6月1日到6月30日的订单先过滤出来否则跨月用户会污染结果。如果目标库是Redshift用Spark写离线脚本算出结果后通过JDBC写入reorder_rate表即可但写入时要用overwrite按分区覆盖防止重复任务把数据翻倍。4.3 留存率用datediff代替递归查询留存率关注的是新用户在第1天、第3天、第7天是否回来。实现思路是先算每个用户的首购日期再把所有日期与首购日期求差按差值日汇总first_buy_date buy_daily.groupBy(userId).agg(min(buy_date).alias(first_date)) retention_daily buy_daily.alias(b) \ .join(first_buy_date.alias(f), userId) \ .selectExpr(f.first_date, datediff(b.buy_date, f.first_date) AS day_gap) \ .filter(day_gap BETWEEN 0 AND 30) \ .groupBy(first_date, day_gap).count()得到的结果表里day_gap0就是首日人数day_gap1是次日留存用户数除以首日人数就得到次日留存率。这个方案把复杂的递归计算变成一次join加groupBy是Spark能高效跑留存的原因。如果留存表的基数很大还可以考虑用Cube预聚合但没必要上来就上先用这个简单方案能跑通再优化。4.4 结果写回时避免幂等问题的两个细节离线任务经常因为重跑覆盖数据如果结果表没有分区或唯一键重复跑会翻倍。写MySQL时我建议先写临时表然后用一条INSERT OVERWRITE/REPLACE INTO写主表常见写法是df.write.format(jdbc) \ .option(truncate, true) \ .option(batchsize, 1000) \ .mode(overwrite) \ .save()但这个overwrite对大规模结果很危险会把整表先truncate再写一旦写入中断结果表就是空的。生产环境建议建分区表把dt分区作为where条件写入模式改成overwrite并且先写到一个临时目录再load data到目标表。另一个细节是batchsize不要超过2000MySQL的max_allowed_packet默认4M一批数据太大容易把数据库连接打挂。这些不是炫技是让你凌晨重跑任务时不用提心吊胆。5. 在YARN上提交Spark作业的实战参数和验证技巧5.1 客户端提交不一定需要完整Spark集群常见问题Spark on YARN提交是不是只要一个Spark客户端答案是的只要提交节点上有Spark发行版、YARN/HDFS客户端配置以及作业依赖的Jar包就能以cluster模式提交。Driver和Executor会在集群的container里启动提交结束后可以关掉终端。整条命令为spark-submit \ --master yarn \ --deploy-mode cluster \ --name ecom_behavior_daily \ --driver-memory 2g \ --executor-memory 6g \ --executor-cores 3 \ --num-executors 20 \ --conf spark.sql.shuffle.partitions240 \ --conf spark.sql.adaptive.enabledtrue \ --class com.ecom.behavior.BehaviorAnalysisApp \ --jars mysql-connector-java-8.0.33.jar \ /data/app/ecom-behavior-analysis_2.12.jar这里的--jars要把MySQL驱动和Redshift驱动都带上否则读写外部数据库时会在executor端抛NoClassDefFoundError。client模式和cluster模式的差别在于Driver运行位置client适合调试cluster适合定时任务。5.2 资源参数表照着调但不硬抄参数建议值说明--num-executors20根据队列资源调整一般每个executor 3-4核是甜点--executor-cores3过大容易导致HDFS I/O竞争--executor-memory6g预留offheap和Overhead设为6-8g为宜--driver-memory2g有collect时调大但尽量避免collect到driverspark.sql.shuffle.partitions200-400核心数的2-3倍不是越大越好这里要算一笔账如果一台机器有32GB内存和8核executor-memory设为6g时一个节点能启动5个executor但你给了20个executor就需要至少4台机器。YARN会根据AM资源自动分配给多了作业会一直Pending。调参时先看yarn node -list确认可用资源再写。5.3 作业验证比看日志更高效的断言技巧跑完Spark作业后用Spark UI看Shuffle spill量和stage耗时能定位严重问题但快速验证结果是否合理更好的办法是在作业末尾加断言。比如算完复购率后加入assert reorder_rate.filter(dt2024-06-01).count() 0 assert 0 reorder_rate.select(reorder_rate).collect()[0][0] 1跑数前先做两个断言比在脚本里print要可靠得多。如果复购率大于1或为负直接抛异常让调度平台标记失败避免脏数据被下游读取。另一个技巧是写个最小化测试取一天的日志在本地跑用pytest断言结果与手工SQL一致再上传集群。本文还有配套的精品资源点击获取
返回列表