
简介这份资源是一套基于Spark技术栈构建的电商用户行为分析大数据平台项目源码面向具备一定Java与大数据基础、希望积累完整项目经验的开发者与学习者。项目围绕用户画像分析、商品推荐算法、实时流量监控、交易数据挖掘与用户行为轨迹追踪等核心模块展开可帮助读者理解Spark在电商场景下的实际应用方式与工程组织思路。压缩包共82个文件以77个Java源码为主体另含pom.xml构建配置、properties参数文件、说明文档与附赠资料整体约138KB结构紧凑便于按模块阅读。目前已有131人学习下载。通过研读源码与配套说明读者可掌握推荐模型计算、实时流量分析、行为路径追踪等功能的实现逻辑并借鉴其目录划分与配置方式为自身项目开发或课程设计提供参考。1. 从一份电商行为日志说起Spark 大型项目实战到底在做什么电商平台每天产生的用户行为日志量级动辄几十 GB 到 TB 级。点一下商品详情、加一次购物车、下一笔订单、搜一个关键词这些动作被埋点 SDK 采集后落到日志文件或消息队列里最终汇入数据仓库。问题在于用传统单机脚本处理这种规模的数据跑一次全量用户画像要几个小时甚至直接 OOM。Spark 大型项目实战电商用户行为分析大数据平台本质上就是用 Spark 技术栈把这套流程工程化——从原始日志清洗、用户画像标签计算、商品推荐算法训练到实时流量监控和交易数据挖掘全部跑在分布式集群上。这套方案适合谁一是正在做大数据课程设计或毕业项目的学生需要一套能跑通、能演示、能写进论文的完整链路二是刚转行大数据的初中级工程师想通过一个真实场景把 Spark SQL、Spark Streaming、MLlib 串起来三是已经在做离线数仓、想补实时链路和推荐算法的从业者。核心诉求就一个别只讲概念告诉我数据怎么进、怎么算、怎么存、怎么查。常见做法是Kafka 接埋点日志Spark Structured Streaming 做实时聚合Hive/ClickHouse 存结果Spark MLlib 跑协同过滤推荐最后用可视化面板展示。下面按这个链路拆开讲。2. 数据接入与清洗从原始日志到可用宽表2.1 日志格式长什么样为什么不能直接拿来算电商埋点日志通常是 JSON 行格式每行一条事件记录。字段包括用户 ID、商品 ID、事件类型view/cart/order/pay、时间戳、页面来源、设备信息、地理区域等。原始日志的问题很集中字段缺失、时间戳格式不统一、事件类型拼写不一致、同一个用户 ID 在不同端可能大小写不同。直接拿这些数据做画像或推荐结果会严重偏移。我一般会先做一轮探查用 Spark SQL 读原始 JSON统计各字段空值率和枚举值分布。这一步不写复杂逻辑就是看清楚数据长什么样。from pyspark.sql import SparkSession from pyspark.sql.functions import col, count, when spark SparkSession.builder \ .appName(EcommerceLogProfiling) \ .config(spark.sql.shuffle.partitions, 200) \ .enableHiveSupport() \ .getOrCreate() # 读取原始 JSON 日志每行一个 JSON 对象 raw_df spark.read.json(hdfs:///data/ecommerce/raw_logs/2024-01-01/) # 统计关键字段空值率 total raw_df.count() null_stats raw_df.select([ (count(when(col(c).isNull(), c)) / total).alias(c _null_rate) for c in [user_id, item_id, event_type, timestamp, price] ]) null_stats.show() # 查看事件类型枚举分布 raw_df.groupBy(event_type).count().orderBy(count, ascendingFalse).show(20)这段代码的逻辑很直接先建 SparkSession开 Hive 支持方便后续写表然后读 JSON 目录Spark 会自动推断 schema。spark.sql.shuffle.partitions设成 200 是经验值集群核数少就降到 50100否则小文件太多反而拖慢。空值率统计用count(when(...))比filter().count()少扫一遍数据。事件类型分布用来发现拼写错误比如 purchase 和 pay 混用。参数说明spark.sql.shuffle.partitions控制 shuffle 后分区数默认 200数据量小于 10GB 时建议设为 CPU 核数的 23 倍。enableHiveSupport()让你能直接saveAsTable写 Hive 表省去手动建 schema。2.2 清洗规则与宽表构建哪些字段必须保留清洗的核心原则是保留分析必需字段统一格式过滤无效记录。具体规则我一般这么定过滤user_id或item_id为空的记录这些无法关联任何维度。时间戳统一转成yyyy-MM-dd HH:mm:ss格式同时拆出dt、hour字段方便分区和实时聚合。事件类型做映射view、cart、order、pay四类其他归为other或直接丢弃。价格字段做异常值处理负数、超过 100 万的置为 null。去重同一用户同一秒对同一商品的重复点击只保留一条。from pyspark.sql.functions import ( from_unixtime, to_timestamp, when, lower, trim, date_format, hour ) from pyspark.sql.types import DoubleType # 清洗与标准化 cleaned_df raw_df \ .filter(col(user_id).isNotNull() col(item_id).isNotNull()) \ .withColumn(user_id, lower(trim(col(user_id)))) \ .withColumn(event_type, when(col(event_type).isin(view, cart, order, pay), col(event_type)) .otherwise(other)) \ .withColumn(event_time, to_timestamp(from_unixtime(col(timestamp)))) \ .withColumn(dt, date_format(col(event_time), yyyy-MM-dd)) \ .withColumn(hour, hour(col(event_time))) \ .withColumn(price, when((col(price).cast(DoubleType()) 0) | (col(price).cast(DoubleType()) 1000000), None) .otherwise(col(price).cast(DoubleType()))) \ .dropDuplicates([user_id, item_id, event_time]) # 写入 Hive 分区表 cleaned_df.write.mode(overwrite) \ .partitionBy(dt) \ .saveAsTable(dwd.ecommerce_user_behavior)逻辑说明lower(trim(...))统一用户 ID 大小写和空格避免同一用户被拆成多个。to_timestamp(from_unixtime(...))处理秒级时间戳如果是毫秒级要去掉后三位。dropDuplicates按用户商品时间去重防止埋点重复上报。最后按dt分区写入后续查询按天过滤能大幅减少扫描量。参数说明partitionBy(dt)按天分区数据量大时可以再加hour做二级分区但分区数太多会导致 NameNode 压力大一般天分区够用。mode(overwrite)适合每日全量重跑增量场景用append。提示清洗后的宽表建议再跑一次空值率和枚举分布确认清洗规则没有误杀。我见过把event_type映射写错导致所有订单被归为 other 的情况血泪经验。3. 用户画像与商品推荐标签计算和协同过滤落地3.1 用户画像标签体系怎么设计才可维护用户画像不是标签越多越好而是标签之间能组合出业务价值。我一般分四层基础属性性别、年龄、地域、行为标签近 7 天浏览量、加购数、下单数、偏好标签偏好品类、价格带、价值标签RFM 分层。每层标签独立计算最后打宽成一张用户宽表。行为标签用 Spark SQL 按时间窗口聚合就行。关键是窗口定义要统一近 7 天、近 30 天、近 90 天别今天用 7 天明天用 10 天。-- 近7天用户行为聚合 INSERT OVERWRITE TABLE dws.user_behavior_7d PARTITION (dt2024-01-07) SELECT user_id, COUNT(CASE WHEN event_type view THEN 1 END) AS view_cnt_7d, COUNT(CASE WHEN event_type cart THEN 1 END) AS cart_cnt_7d, COUNT(CASE WHEN event_type order THEN 1 END) AS order_cnt_7d, COUNT(CASE WHEN event_type pay THEN 1 END) AS pay_cnt_7d, SUM(CASE WHEN event_type pay THEN price ELSE 0 END) AS pay_amount_7d, COLLECT_SET(CASE WHEN event_type IN (cart,pay) THEN item_id END) AS interacted_items FROM dwd.ecommerce_user_behavior WHERE dt BETWEEN 2024-01-01 AND 2024-01-07 GROUP BY user_id;逻辑说明COUNT(CASE WHEN ...)一次扫描算出多个指标比多次filter().count()高效。COLLECT_SET收集用户交互过的商品 ID去重后用于推荐算法的输入。pay_amount_7d是消费金额用于 RFM 中的 MMonetary。参数说明dt BETWEEN范围根据窗口调整近 30 天就改日期范围。COLLECT_SET在用户交互商品多时可能产生大数组建议加LIMIT或改用APPROX_COUNT_DISTINCT只算数量。偏好标签需要关联商品维度表按品类和价格带聚合。价值标签用 RFM 模型R最近一次消费距今天数、F消费频次、M消费金额每项分 5 档组合成 125 种分层。实际业务中通常简化为 8 类核心人群。3.2 协同过滤推荐ALS 参数怎么调才不翻车商品推荐用 Spark MLlib 的 ALS交替最小二乘做协同过滤输入是用户-商品-评分矩阵。电商场景没有显式评分用行为加权构造隐式反馈view1cart3order5pay10。然后调 ALS 训练。from pyspark.ml.recommendation import ALS from pyspark.ml.evaluation import RegressionEvaluator from pyspark.sql.functions import col, when # 构造隐式评分 ratings cleaned_df.select(user_id, item_id, event_type) \ .withColumn(rating, when(col(event_type) view, 1.0) .when(col(event_type) cart, 3.0) .when(col(event_type) order, 5.0) .when(col(event_type) pay, 10.0) .otherwise(0.0)) \ .filter(col(rating) 0) \ .groupBy(user_id, item_id).agg({rating: max}) \ .withColumnRenamed(max(rating), rating) # 划分训练集和测试集 (training, test) ratings.randomSplit([0.8, 0.2], seed42) # ALS 训练 als ALS( maxIter10, regParam0.01, userColuser_id, itemColitem_id, ratingColrating, implicitPrefsTrue, coldStartStrategydrop, nonnegativeTrue ) model als.fit(training) # 评估 predictions model.transform(test) evaluator RegressionEvaluator(metricNamermse, labelColrating, predictionColprediction) rmse evaluator.evaluate(predictions) print(fRMSE: {rmse}) # 为每个用户生成 Top10 推荐 user_recs model.recommendForAllUsers(10) user_recs.show(5, truncateFalse)逻辑说明groupBy(user_id, item_id).agg({rating: max})取用户对同一商品的最大行为权重避免重复行为叠加。implicitPrefsTrue表示隐式反馈ALS 会把评分当作置信度处理。coldStartStrategydrop丢弃测试集中出现但训练集中没有的用户/商品否则预测值为 NaN 导致 RMSE 计算失败。nonnegativeTrue约束因子非负推荐结果更稳定。参数说明maxIter一般 1020再多收益递减。regParam正则化系数0.010.1 之间调过拟合就加大。rank因子维度默认 10数据量大可以调到 50100但训练时间和内存同步上升。RMSE 在隐式反馈场景下参考价值有限更该看召回率K 或人工抽查推荐结果。注意ALS 训练前一定要检查用户和商品 ID 是否已转成整数索引。MLlib 要求userCol和itemCol是 int 类型字符串 ID 需要先StringIndexer编码否则直接报错。4. 实时流量监控与交易数据挖掘Structured Streaming 链路4.1 实时流量监控窗口聚合和延迟处理实时流量监控要回答的问题很具体当前每分钟 PV/UV 多少、哪个商品被访问最多、下单转化率有没有异常波动。用 Spark Structured Streaming 从 Kafka 消费日志做 1 分钟滚动窗口聚合结果写 ClickHouse 或 Redis 供面板查询。from pyspark.sql.functions import window, countDistinct, count, from_json from pyspark.sql.types import StructType, StringType, LongType, DoubleType # 定义 Kafka 消息中 JSON 的 schema schema StructType() \ .add(user_id, StringType()) \ .add(item_id, StringType()) \ .add(event_type, StringType()) \ .add(timestamp, LongType()) \ .add(price, DoubleType()) # 从 Kafka 读取 stream_df spark.readStream \ .format(kafka) \ .option(kafka.bootstrap.servers, kafka-broker:9092) \ .option(subscribe, ecommerce_behavior) \ .option(startingOffsets, latest) \ .load() # 解析 JSON 并做窗口聚合 parsed_df stream_df.select( from_json(col(value).cast(string), schema).alias(data) ).select(data.*) windowed parsed_df \ .withColumn(event_time, to_timestamp(from_unixtime(col(timestamp)))) \ .withWatermark(event_time, 2 minutes) \ .groupBy(window(col(event_time), 1 minute), col(item_id)) \ .agg( count(*).alias(pv), countDistinct(user_id).alias(uv) ) # 输出到 ClickHouse query windowed.writeStream \ .outputMode(update) \ .foreachBatch(lambda df, epoch_id: df.write \ .format(jdbc) \ .option(url, jdbc:clickhouse://clickhouse:8123/default) \ .option(dbtable, realtime_traffic) \ .option(user, default) \ .mode(append) \ .save()) \ .trigger(processingTime30 seconds) \ .start() query.awaitTermination()逻辑说明from_json按 schema 解析 Kafka value避免手动拆字段。withWatermark(event_time, 2 minutes)允许 2 分钟延迟数据超过则丢弃防止状态无限增长。window(col(event_time), 1 minute)做 1 分钟滚动窗口。outputMode(update)只输出有更新的窗口减少下游写入压力。foreachBatch里写 ClickHouse用 JDBC 批量插入。参数说明startingOffsets设latest只消费新数据重跑历史用earliest。trigger(processingTime30 seconds)每 30 秒触发一次微批延迟和吞吐的折中。watermark设 2 分钟是根据业务容忍延迟定的要求更低延迟就设 30 秒但乱序数据会被丢弃更多。4.2 交易数据挖掘从订单表里挖出异常和规律交易数据挖掘常见两个方向异常订单检测和商品关联规则。异常订单用统计方法就行——同一用户短时间大量下单、收货地址频繁变更、订单金额偏离该用户历史均值 3 倍以上。用 Spark SQL 窗口函数算。-- 检测同一用户1小时内下单超过5笔的异常 SELECT user_id, order_time, order_cnt_1h, total_amount_1h FROM ( SELECT user_id, order_time, COUNT(*) OVER (PARTITION BY user_id ORDER BY order_time RANGE BETWEEN INTERVAL 1 HOUR PRECEDING AND CURRENT ROW) AS order_cnt_1h, SUM(amount) OVER (PARTITION BY user_id ORDER BY order_time RANGE BETWEEN INTERVAL 1 HOUR PRECEDING AND CURRENT ROW) AS total_amount_1h FROM dwd.order_detail WHERE dt 2024-01-07 ) t WHERE order_cnt_1h 5 ORDER BY order_cnt_1h DESC;逻辑说明RANGE BETWEEN INTERVAL 1 HOUR PRECEDING AND CURRENT ROW按时间范围滑动比ROWS更符合业务语义。order_cnt_1h 5是阈值实际要根据业务调整大促期间阈值可以放宽。商品关联规则用 FP-Growth 算法输入是每个订单的商品集合。from pyspark.ml.fpm import FPGrowth # 按订单聚合商品列表 basket cleaned_df.filter(col(event_type) pay) \ .groupBy(order_id) \ .agg(collect_set(item_id).alias(items)) \ .filter(size(col(items)) 1) # FP-Growth 挖掘频繁项集 fp FPGrowth(itemsColitems, minSupport0.01, minConfidence0.3) model fp.fit(basket) # 查看频繁项集和关联规则 model.freqItemsets.show(10, truncateFalse) model.associationRules.show(10, truncateFalse)逻辑说明collect_set(item_id)把同一订单的商品聚成数组size 1过滤单商品订单。minSupport0.01表示商品组合至少出现在 1% 的订单中数据量大可以降到 0.001。minConfidence0.3表示规则置信度至少 30%。参数说明minSupport越低挖掘出的规则越多但噪声也越大建议从 0.01 开始试。minConfidence根据业务对准确率的要求调推荐场景一般 0.20.5。5. 避坑与排查这套链路最容易翻车的 5 个地方5.1 数据倾斜导致任务卡在 99%现象Spark 任务跑到 99% 卡住几小时看 UI 发现某个 task 处理的数据量是其他的几十倍。原因某个热门商品或测试用户产生了海量日志groupBy(item_id)时全分到一个分区。解决对倾斜 key 加随机前缀打散聚合后再合并或者开启spark.sql.adaptive.enabledtrue让 AQE 自动处理倾斜 join。5.2 ALS 训练报 NaN 或 RMSE 为 NaN现象模型训练完predictions里 prediction 全是 NaNRMSE 也是 NaN。原因测试集里有训练集没出现过的用户或商品ALS 无法预测。解决设coldStartStrategydrop丢弃这些记录或者用StringIndexer时设handleInvalidskip过滤未知类别。5.3 Structured Streaming 状态无限增长现象实时任务跑几天后内存暴涨checkpoint 目录越来越大。原因没设 watermark 或 watermark 设得太宽松窗口状态一直保留。解决必须设withWatermark时间设成业务最大容忍延迟定期清理 checkpoint 目录用outputMode(update)而不是complete。5.4 小文件过多拖垮 NameNode现象清洗后写 Hive 表每个分区几千个小文件后续查询启动几百个 task。原因spark.sql.shuffle.partitions设太大或者流式写入频繁触发小批。解决离线写入前用repartition(200)控制文件数流式场景用trigger拉长微批间隔定期跑ALTER TABLE ... CONCATENATE合并小文件。5.5 时间戳时区不一致导致窗口错位现象实时监控面板显示的数据比实际晚 8 小时。原因埋点时间戳是 UTCSpark 默认时区是 UTC但业务看的是北京时间。解决spark.conf.set(spark.sql.session.timeZone, Asia/Shanghai)或者在from_unixtime后加 8 小时偏移。统一在清洗层处理别在展示层补。6. 进阶技巧用广播变量和缓存把重复计算压下去这套链路跑通之后性能瓶颈通常不在算法而在重复扫描和 shuffle。我自己的习惯是任何被多次引用的 DataFrame先cache()再操作小维度表商品品类、地域映射用broadcast()广播避免 join 时 shuffle。from pyspark.sql.functions import broadcast # 商品维度表通常只有几万行广播后 join 不产生 shuffle item_dim spark.table(dim.item_info) behavior spark.table(dwd.ecommerce_user_behavior) enriched behavior.join(broadcast(item_dim), onitem_id, howleft) enriched.cache() enriched.count() # 触发缓存 # 后续多次使用 enriched 做不同聚合不会重复读 HDFS enriched.groupBy(category).count().show() enriched.groupBy(brand).agg({price: avg}).show()逻辑说明broadcast(item_dim)把小表分发到每个 executorjoin 时在本地完成省掉 shuffle。cache()把 enriched 缓存在内存后续两次groupBy直接读缓存。count()是触发缓存的 action不调用的话 cache 不会生效。参数说明广播表默认阈值 10MB超过会报错可以调spark.sql.autoBroadcastJoinThreshold但别超过 200MB否则 executor 内存扛不住。cache()用 MEMORY_AND_DISK 级别更稳内存不够会溢写到磁盘。验证缓存是否生效看 Spark UI 的 Storage 页面确认 RDD 的 Fraction Cached 是 100%。如果多次查询时间没变化大概率是缓存没触发或者被 evict 了。另一个技巧是分区裁剪。清洗后的表按dt分区查询时一定带上WHERE dt ...否则全表扫描。我见过有人写WHERE dt 2024-01-01结果扫了整年数据跑了一整夜。最后说一个我踩过的坑Spark 3.x 默认开启 AQE但spark.sql.adaptive.coalescePartitions.enabled在某些版本有 bug导致分区合并后数据不均。如果发现任务并行度突然下降先把这个参数关掉试试。这些参数没有银弹都是根据集群规模和数据特征一点点调出来的。希望帮到你。本文还有配套的精品资源点击获取