简介:基于Spark技术栈的电商用户行为分析系统实战项目,面向大数据开发与电商数据分析学习者,覆盖用户画像、商品推荐、实时流量监控、交易数据挖掘与行为轨迹追踪等核心场景,演示在一个统一工程中整合多种分析逻辑的实现方式。包内共82个文件,以77个Java源码文件为核心,负责各功能模块的业务逻辑;XML与Properties提供配置,MD、TXT、DOCX文档包含项目说明与操作指引;压缩包仅138KB,轻量易用。目前已有132人浏览学习,适合希望从代码层面掌握项目落地方式的数据开发初中级学习者。将项目导入IDE后,可以对照文档快速定位画像、推荐、监控、挖掘、轨迹等模块的代码路径,梳理其数据结构与调用关系;整体工程不大,却提供了完整的模块划分和工程组织范例,适合课程设计、毕业设计,也可作为企业级电商行为分析平台二次开发的参考。
1. 电商用户行为分析大数据平台:为什么我建议你直接上 Spark
电商大促当晚,点击、浏览、加购、下单、支付、售后这条链路每分钟都会产生上亿条行为日志,MySQL 里那几张订单表根本接不住这种量级。Spark 大型项目实战里最常见也最成熟的方向,就是用一套 Spark 技术栈把用户行为日志接进来,统一做离线画像、商品推荐、实时流量监控和交易数据挖掘。这个电商用户行为分析大数据平台,本质是把批计算和流计算放在同一条技术线上,让数仓、算法、运营看同一份数据。适合正在搭数据平台的数据开发、数仓工程师,也适合想从业务报表往用户画像和推荐方向转的团队。下文我会按实际落地方案讲架构、代码和参数,并把容易翻车的地方单独拉出来。
2. 先定架构和表模型:平台不散,后面才不返工
做电商用户行为分析最怕的不是规模大,而是每个人对“用户行为”理解不一样。产品说点击,运营说页面停留,算法说埋点序列,最后各算各的,数据对不上。所以第一件事不是写 Spark 作业,而是把整体分层和数据模型定死。
2.1 分层设计:离线数仓与实时链路共用一套计算层
我一般把平台分成五层。ODS 层直接接原始日志,HDFS 上按天分区的 ORC 文件;DWD 层做清洗、去重、会话划分,输出标准行为事件表;DWS 层做轻聚合,比如用户日维度的浏览、加购、下单次数;ADS 层直接服务 BI 报表和画像接口。实时只有一条链路:Kafka 进 Spark Structured Streaming,做分钟级窗口聚合,结果落到 Redis 和 ES,不对接 Hive 数仓。
计算层统一用 Spark,批跑 Spark SQL,流跑 Structured Streaming,资源由 Yarn 统一调度。我的参数起点是每个 executor 分配 4 GB 内存、2 核,spark.sql.shuffle.partitions 设 200。数据量上来后,这两个参数往往是第一个要调的,后面避坑章细说。存储层用 HDFS 做离线主存储,HBase 存画像明细,Redis 存推荐结果和实时指标,ES 只承担行为轨迹检索。
2.2 行为埋点表模型:一张宽表撑起全部业务
行为日志的最佳落地方式是一张事件宽表,字段统一,禁止各业务线各自建表。核心字段有:event_id、user_id、product_id、session_id、event_type、event_time、page_id、device_type、channel,以及一个放扩展属性的 map 字段,用于承接不同业务的埋点参数。
CREATE TABLE IF NOT EXISTS dwd.dwd_user_behavior_event ( event_id STRING COMMENT '事件唯一ID', user_id STRING COMMENT '用户ID,登录用户', visitor_id STRING COMMENT '匿名访客ID,未登录前用', product_id STRING COMMENT '商品ID,无商品行为则为空', session_id STRING COMMENT '会话ID,由代码统一生成', event_type STRING COMMENT 'view/cart/order/pay/favorite', event_time BIGINT COMMENT '事件时间,毫秒时间戳', page_id STRING COMMENT '页面ID', device_type STRING COMMENT 'android/ios/pc/h5', channel STRING COMMENT '渠道标识', ext_map MAP<STRING,STRING> COMMENT '扩展属性' ) PARTITIONED BY (dt STRING) STORED AS ORC;这张表设计的关键是 event_time 用毫秒级 BIGINT,而不是字符串。字符串时间做范围过滤和 session 切分时,转格式的代价很高。session_id 也不要靠埋点上报,客户端生成的 session 值经常因为 WebView 重载而断掉,最佳做法是在 DWD 层按规则重新切分。ORC 加列存压缩能省一半以上存储,日常分析也只取需要的列。
2.3 会话切分口径:同一个用户两条轨迹怎么拼
行为轨迹还原的前提是会话切分一致。常见口径是 30 分钟内无新行为则会话结束,下一个动作属于新会话。这个逻辑在 Spark SQL 里用 LAG 函数就能做,注意要对 user_id 分区排序后再算相邻事件的时间差。另外要统一 user_id 与 visitor_id:登录用户直接用 user_id,未登录用户先用 visitor_id 聚合,登录后再做一次映射归并。
时间口径也要收紧。埋点客户端上报的时间和服务器接收时间往往有偏差,移动端还可能因为本地时钟错误产生未来时间。我在 DWD 层统一以服务器接收时间为准,字段名就叫 event_time,客户端上报的原始时间存进 ext_map 里只做参考。这个决定能避免后来做漏斗分析时数据“倒挂”,比如支付时间早于下单时间。
3. 离线用户画像与轨迹还原:Spark 读 JSON 和会话切分的完整写法
画像和轨迹是离线分析的核心。先讲怎么读原始日志,再讲怎么切会话拼轨迹,最后讲 RFM 画像怎么落库。
3.1 读取 JSON 行为日志:先定 Schema,再交给 Spark
很多团队第一步就吃 read.json 自动推断类型的亏。行为日志里 ext_map 是动态结构,自动推断会把所有整数推断成 BIGINT,把时间戳推断成 STRING,等你要 join 的时候就傻眼。我的做法是先手动定义 Schema。
from pyspark.sql.types import ( StructType, StructField, StringType, LongType, MapType ) behavior_schema = StructType([ StructField("event_id", StringType()), StructField("user_id", StringType()), StructField("visitor_id", StringType()), StructField("product_id", StringType()), StructField("event_type", StringType()), StructField("ts", LongType()), # 客户端时间戳 StructField("server_time", LongType()), # 服务器接收时间戳 StructField("device_type", StringType()), StructField("channel", StringType()), StructField("ext_map", MapType(StringType(), StringType())) ]) df = spark.read \ .option("multiLine", False) \ .schema(behavior_schema) \ .json("hdfs:///raw/behavior_log/dt=2024-11-11/") df = df.filter(df["event_type"].isin( "view", "cart", "order", "pay", "favorite" ))指定 Schema 之后,Spark 不会再去扫描一遍文件推断类型,读取速度明显更快。过滤事件类型这步放在读取后立即做,早过滤早减少下游 shuffle 数据量。这里我没有用 dropDuplicates,因为行为日志天然有重复的可能,去重要放到会话切分后按 event_id 做,避免把两个独立会话合并掉。
3.2 会话切分与行为轨迹拼接:LAG 算间隔,COLLECT_LIST 拼链路
会话切分在 DWD 层做一次,全平台共用。按 user_id 分区,按 server_time 排序,计算每条记录与上一条的时间差,超过 30 分钟就生成新的 session_id。
from pyspark.sql.window import Window from pyspark.sql.functions import ( lag, when, sum, col, concat_ws, collect_list ) w = Window.partitionBy("user_id").orderBy("server_time") df = df.withColumn( "prev_time", lag("server_time", 1).over(w) ).withColumn( "time_gap", when(col("prev_time").isNotNull(), col("server_time") - col("prev_time")) .otherwise(0) ).withColumn( "is_new_session", when(col("time_gap") >= 1800000, 1).otherwise(0) ).withColumn( "session_id", concat_ws("_", "user_id", sum("is_new_session").over(w)) ) track_df = df.groupBy("user_id", "session_id") \ .agg(collect_list( concat_ws(":", "event_type", "product_id", "server_time") ).alias("track"))这里的窗口从 user_id 分区,意味着同一用户的所有事件都会分到同一个 executor。用户量如果极不均匀,大用户会拖慢整个 Stage,所以我一般会在前面加一步 repartition(user_id, 500),让 Spark 按哈希分散到 500 个分区再开窗口。轨迹字段用 event_type:product_id:server_time 拼接,后面做路径分析时用 split 拆开即可,不需要一开始就设计嵌套结构。
3.3 画像特征落 HBase:RFM 不只是算三个数
画像里最常用的是 RFM。R 看最后一次下单距今多少天,F 看近 90 天下单次数,M 看近 90 天支付金额。但直接对全量用户做聚合很费资源,一般只对近 30 天有活跃行为的用户计算,老沉睡用户不做更新。
from pyspark.sql.functions import datediff, current_date, count, sum rfm_df = df.filter( (col("event_type") == "pay") & (col("server_time") >= start_of_90d) ).groupBy("user_id").agg( datediff(current_date(), max(col("server_time"))).alias("recency"), count("*").alias("frequency"), sum(col("ext_map")["pay_amount"].cast("double")).alias("monetary") ) def label_rfm(recency, frequency, monetary): if monetary >= 500 and frequency >= 5: return "high_value" elif monetary >= 100: return "medium_value" return "low_value" rfm_df.createOrReplaceTempView("rfm_tmp") result = spark.sql(""" SELECT user_id, CASE WHEN monetary >= 500 AND frequency >= 5 THEN 'high_value' WHEN monetary >= 100 THEN 'medium_value' ELSE 'low_value' END AS rfm_level FROM rfm_tmp """)这里要留意 pay_amount 存在 ext_map 里,它在 JSON 中本质是字符串,必须先 cast 成 double 再做 sum,否则结果全是 0 或者报类型错误。画像结果落地 HBase 用 bulkput,rowkey 用 user_id 倒序,避免顺讯写入全部打到同一个 Region。
4. 实时流量监控、交易挖掘与商品推荐:一条流三条出口
实时这块真正难的不是窗口统计,而是把实时指标和离线指标口径对齐。我踩过最狠的坑是:实时 UV 用 Redis Set 去重,离线 UV 用 count(distinct user_id),两边永远对不上。后来统一用“日首次活跃”规则,实时与离线都按当天第一次出现的 user_id 计一次,数字才对得上。
4.1 Structured Streaming 消费 Kafka:窗口聚合与去重计数
实时流量监控我用 Spark Structured Streaming 消费 Kafka 的 dwd_user_behavior_event 主题。数据源已经过 DWD 清洗,流任务这边只做窗口聚合。
from pyspark.sql.functions import window, approx_count_distinct kafka_df = spark.readStream \ .format("kafka") \ .option("kafka.bootstrap.servers", "kafka1:9092,kafka2:9092") \ .option("subscribe", "dwd_user_behavior_event") \ .option("startingOffsets", "latest") \ .load() behavior_df = kafka_df.selectExpr( "cast(value as string) as json_str" ).select( from_json("json_str", behavior_schema).alias("data") ).select("data.*") stream_query = behavior_df \ .withWatermark("server_time", "10 minutes") \ .groupBy( window("server_time", "5 minutes", "1 minute"), "channel" ).agg( count("*").alias("pv"), approx_count_distinct("user_id").alias("uv") ) \ .writeStream \ .outputMode("update") \ .format("console") \ .option("truncate", "false") \ .trigger(processingTime="60 seconds") \ .start()withWatermark 给 10 分钟,是业务上对“迟到数据”的最大容忍度。窗口用了 5 分钟滑动、1 分钟步长,指标每分钟刷新一次,但近 5 分钟区间内不会因单点抖动大起大落。approx_count_distinct 是近似去重,误差一般在 1% 以内,UV 这种量级完全够用。如果严格要准确去重,换 count(distinct) 但内存压力会成倍上涨。
4.2 交易数据挖掘:RFM 分层与复购周期不是每天重跑
交易挖掘我并不会每天都全量重算,代价太高。常见做法是“首单初始化 + 每日增量更新”。每天只算当天有订单行为的用户,把频次、金额累加到存量画像表;30 天没有新行为的用户冻结画像,不再参与每日增量。这样 ADS 层的数据量被压得很小,跑一遍也就十分钟。
复购周期这个指标看着简单,写起来容易错。比如用户第 1 天下单,第 3 天下单,第 10 天下单,那他有 2 个“间隔”。直接算平均间隔,会把第 1 天到第 10 天这 9 天也算进去,时间跨度比实际复购周期偏大一截。正确做法是先对每个用户的下单时间排序,用 LEAD 函数取下一次下单时间,再算差值。
4.3 商品推荐离线召回:协同过滤 ALS 的参数经验
推荐部分主体用 ALS 做离线召回,它在 Spark MLlib 里开箱即用,适合做用户–商品行为矩阵的隐式反馈分解。数据构造我用了分行为权重:点击算 1,加购算 3,下单算 5,支付算 8,支付金额作为附加权重,没有真实评分也照样训练。
from pyspark.ml.recommendation import ALS train_df = behavior_df.filter( col("event_type").isin("view", "cart", "order", "pay") ).groupBy("user_id", "product_id").agg( sum(when(col("event_type") == "view", 1) .when(col("event_type") == "cart", 3) .when(col("event_type") == "order", 5) .otherwise(8)).alias("rating") ) als = ALS( userCol="user_id", itemCol="product_id", ratingCol="rating", implicitPrefs=True, rank=50, regParam=0.01, alpha=40, maxIter=10, coldStartStrategy="drop" ) model = als.fit(train_df) user_recs = model.recommendForAllUsers(20)implicitPrefs 必须设成 True,因为我们的 rating 是隐式反馈而非显式评分。alpha 控制隐式反馈的置信度,取 40 是业界常见起点,调大可让“只点不买”的用户更偏向于已购商品。rank 先给 50,数据量到千万级后加到 100。coldStartStrategy 设 drop,是避免训练集里没有的新用户在预测时报 null。
5. 常见问题与避坑指南:这五个坑足够让平台重跑三天
这一章全是实战里真金白银换来的教训。每一条都按现象、原因、解决三步写,照着排查能省一天时间。
5.1 资源与调度:两类最隐蔽的翻车现场
第一个坑:Spark UI 显示每个 Task 处理的数据量均匀,但总时长被少数几个 Task 拖住,整个 Stage 跑几十分钟。原因是按 user_id 或 product_id 做 groupBy 时,头部用户和爆款商品产生的同 key 数据量太大,哈希分区解决不了这种数据倾斜。解决方法是先做随机前缀打散,两步聚合,或者用 salted key 分桶后再 join。我用得最多的是先把热点 key 单独拎出来,用 repartition 按子 key 分散,效果比调大 shuffle 分区更直接。
第二个坑:Structured Streaming 任务改了消费逻辑,重启时报 checkpoint 里已有旧的 committed offset,新增字段后恢复位置总报错。原因是同一个 checkpoint 路径里保存了旧的执行计划元数据,代码一变,元数据对不上。解决方法是永远不要复用同一个 checkpoint 跑两份不同逻辑的流;开发环境每改一次 schema,就换一个新的 checkpoint 路径。生产上要升级逻辑,先停流,备份 checkpoint 目录,再指向新目录启动,确认数据没有重复消费后再删旧目录。这里没有后悔药,路径一旦切错就丢数据。
5.2 数据质量:结果不对,比跑不动更可怕
第三个坑:收入日报的支付金额总和比财务后台少了几十万。原因是支付的扩展属性里,部分金额字段在旧版本客户端上报的是字符串“12.5”,新版本上报的是数字 12.5,Spark 读取时统一按字符串处理,但老数据的值在解析时被转成 null。解决方法是定义 Schema 时把所有金额字段都声明成字符串,在 DWD 层统一用 cast 清洗成 decimal,同时过滤掉“金额为 0 且事件为 pay”的异常记录。数据管道对上游格式的假设越少,越经得起上层改版。
第四个坑:会话切分后,轨迹里出现“下单 → 浏览 → 支付”这种倒序行为。原因是 session 切分时混用了客户端时间 ts,而客户端时间存在本地时钟偏移。解决方法是整个平台只在 DWD 层使用 server_time,并把 ts 放进扩展属性仅做数据质量监控;如果发现同一会话里超过 50% 的事件是乱序,说明这条链路的服务器接收时间本身有问题,要去查 NTP 同步。
第五个坑:ALS 推荐结果里大量商品是全局热门,完全看不出个性。原因是隐式反馈矩阵太稀疏,用户大多只点过一两个商品,模型等于在用全局均值拟合。解决方法是先做商品维度过滤,去掉 7 天内无任何行为的商品,再对单用户行为次数做对数变换,缓解头部用户权重失衡。如果换成 rank 50 还是热门兜底,就加一层 rule 过滤:推荐结果里同品类商品最多占一半,剩下的位置给长尾商品。
6. 上线前的验证与调优:用一轮“比数”把平台钉死
平台上线前,我不跑冒烟测试,直接做一轮全量对账:拿同一天的原始日志,分别用 Spark 和 Hive SQL 算一遍核心指标,总数、UV、GMV、漏斗转化率,逐项比对。误差在 0.5% 以内放行,超过就要查原因。这个习惯帮我挡下了至少三次字段类型变更引发的线上事故。
验证脚本我收集成一个 run_validate.sh,核心逻辑很简单:先把 Spark 作业产出的结果表导出成临时表,再用 Hive 算同一指标。两边对不上时,先在公共维度找差异源,比如分区缺失、时间口径、未登录用户的怪值。数据量和分区规范都确认后,最后调 Spark 性能参数,不要反过来。
调优时不要盲改并行度。我一般先开 Spark UI 看两个地方:一个是每个 Stage 的 shuffle read 总量,过大就把 shuffle partitions 调高;另一个是某个 Task 的 duration 明显高于同 Stage 其他 Task,先查数据倾斜,而不是加内存。executor 内存我很少超过 8 GB,超过就怀疑是分区策略有问题。
还有一个习惯坚持了很久:每次改动 DWD 层逻辑,都在配置中心留一条版本记录,注明改了什么口径、谁改的、影响哪些下游任务。平台跑久了,真正贵的不是算力,而是“为什么数字会变”的答疑时间。早期我总以为跑通 spark-submit 就算成功,后来发现能说清楚每个数字怎么来的,才算真上线。希望帮到你。
本文还有配套的精品资源,点击获取