【免费】基于Spark实时电商用户行为分析与预测 系统(Python版本+pyspark+可视化大屏+Kafka+FastAPI+Vue3) 锋哥原创出品,必属精品
大家好,我是Java1234_小锋老师,分享一套锋哥原创的基于Spark实时电商用户行为分析与预测 系统(Python版本+pyspark+可视化大屏+Kafka+FastAPI+Vue3)
项目介绍
随着电子商务规模持续扩大,用户在浏览、加购、收藏与购买等环节产生的行为数据呈现高并发、高吞吐与强时效特征。传统离线批处理分析难以满足运营决策对实时性的要求。本文设计并实现了一套基于 Spark 的实时电商用户行为分析与预测系统,围绕“数据采集—流式计算—指标落库—可视化展示—销售预测”的完整链路展开研究与工程实践。
系统采用前后端分离架构:前端基于 Vue3、Element Plus 与 ECharts 构建管理端与数据大屏;后端采用 Python FastAPI 提供 RESTful 接口,并结合 JWT 完成管理员身份认证;实时链路以 Kafka 作为消息中间件承接行为事件,以 Spark Structured Streaming 完成按小时窗口的 PV、UV、加购、收藏、购买与销售额聚合;预测模块基于 Spark ML 线性回归对销售额序列进行建模,并输出 RMSE、MAE、MAPE 等误差指标。数据持久化采用 MySQL,数据库名为 db_ecommerce,核心业务表均以 t_ 前缀命名。
测试结果表明,系统能够稳定完成管理员登录、个人中心维护、行为与商品管理、实时统计展示、销售预测对比及流水线状态监控等功能,具备较好的可扩展性与教学示范价值,可为电商运营提供实时洞察与辅助决策支持。
本文的主要工作包括:完成系统需求分析与总体架构设计;绘制实体属性图与实体关系图并完成八张核心业务表设计;实现基于 Kafka 与 Spark 的实时统计及销售预测链路;完成 Vue3 管理端与数据大屏;开展功能测试并给出改进方向。研究结果表明,将流式计算与 Web 管理系统结合,能够在本科毕业设计条件下形成完整、可运行、可解释的实时分析应用。
源码下载
链接: https://pan.baidu.com/s/1u0yzt7SCx13nEjaH4FSMLw?pwd=1234
提取码: 1234
系统展示
![]()
![]()
![]()
![]()
核心代码
""" Spark ML 销售额预测模块 """ import numpy as np from decimal import Decimal from config import settings def compute_error_metrics(y_true: list, y_pred: list) -> dict: """ 计算误差指标:RMSE、MAE、MAPE """ y_true = np.array(y_true, dtype=float) y_pred = np.array(y_pred, dtype=float) rmse = float(np.sqrt(np.mean((y_true - y_pred) ** 2))) mae = float(np.mean(np.abs(y_true - y_pred))) mask = y_true != 0 if mask.any(): mape = float(np.mean(np.abs((y_true[mask] - y_pred[mask]) / y_true[mask])) * 100) else: mape = 0.0 return {"rmse": round(rmse, 4), "mae": round(mae, 4), "mape": round(mape, 4)} def run_spark_prediction(sales_series: list = None) -> tuple: """ 使用 Spark ML 进行销售额预测 返回 (predictions, error_metrics) """ try: from pyspark.sql import SparkSession from pyspark.ml.feature import VectorAssembler from pyspark.ml.regression import LinearRegression from pyspark.sql.types import StructType, StructField, DoubleType, IntegerType, StringType import pyspark.sql.functions as F spark = SparkSession.builder \ .appName("SalesPrediction") \ .master(settings.SPARK_MASTER) \ .config("spark.driver.memory", "2g") \ .getOrCreate() spark.sparkContext.setLogLevel("WARN") if sales_series is None: from database import SessionLocal from models.realtime_stat import RealtimeStat db = SessionLocal() stats = db.query(RealtimeStat).order_by(RealtimeStat.window_time.asc()).all() db.close() sales_series = [ {"window_time": s.window_time, "sales": float(s.sales)} for s in stats ] if len(sales_series) < 5: spark.stop() return [], {"rmse": 0, "mae": 0, "mape": 0} # 构造滞后特征 data = [] for i in range(3, len(sales_series)): data.append({ "window_time": sales_series[i]["window_time"], "lag1": sales_series[i - 1]["sales"], "lag2": sales_series[i - 2]["sales"], "lag3": sales_series[i - 3]["sales"], "hour": int(str(sales_series[i]["window_time"])[11:13]), "sales": sales_series[i]["sales"], }) schema = StructType([ StructField("window_time", StringType()), StructField("lag1", DoubleType()), StructField("lag2", DoubleType()), StructField("lag3", DoubleType()), StructField("hour", IntegerType()), StructField("sales", DoubleType()), ]) df = spark.createDataFrame(data, schema) # 划分训练集和测试集(后20%作为测试) split_idx = max(int(len(data) * 0.8), 1) train_df = df.limit(split_idx) test_df = df.filter(F.monotonically_increasing_id() >= split_idx) assembler = VectorAssembler( inputCols=["lag1", "lag2", "lag3", "hour"], outputCol="features" ) train_df = assembler.transform(train_df) test_df = assembler.transform(test_df) lr = LinearRegression(featuresCol="features", labelCol="sales", maxIter=100) model = lr.fit(train_df) predictions_raw = model.transform(test_df).collect() predictions = [] y_true, y_pred = [], [] for row in predictions_raw: true_val = float(row["sales"]) pred_val = float(row["prediction"]) predictions.append({ "window_time": row["window_time"], "true_sales": round(true_val, 2), "pred_sales": round(max(pred_val, 0), 2), }) y_true.append(true_val) y_pred.append(pred_val) error = compute_error_metrics(y_true, y_pred) spark.stop() return predictions, error except Exception as e: print(f"[Spark ML] 预测失败: {e}") return None, None def save_predictions_to_db(predictions: list, error: dict): """ 保存预测结果和误差指标到数据库 """ from database import SessionLocal from models.prediction import Prediction from models.error_metric import ErrorMetric db = SessionLocal() try: # 清空旧预测数据 db.query(Prediction).delete() for p in predictions: db.add(Prediction( window_time=p["window_time"], true_sales=Decimal(str(p["true_sales"])), pred_sales=Decimal(str(p["pred_sales"])), )) db.add(ErrorMetric( rmse=Decimal(str(error["rmse"])), mae=Decimal(str(error["mae"])), mape=Decimal(str(error["mape"])), )) db.commit() print(f"[Spark ML] 已保存 {len(predictions)} 条预测结果") finally: db.close()<template> <div class="page-container"> <div class="page-card"> <div class="page-title">销售额预测分析</div> <!-- 误差指标卡片 --> <div class="error-cards"> <div class="error-card"> <div class="metric-label">RMSE (均方根误差)</div> <div class="metric-value">{{ errorMetric.rmse }}</div> </div> <div class="error-card"> <div class="metric-label">MAE (平均绝对误差)</div> <div class="metric-value">{{ errorMetric.mae }}</div> </div> <div class="error-card"> <div class="metric-label">MAPE (平均绝对百分比误差 %)</div> <div class="metric-value">{{ errorMetric.mape }}%</div> </div> </div> <!-- 真实 vs 预测对比图 --> <div ref="compareRef" class="pred-chart pred-chart-compare"></div> <!-- 残差图 --> <div ref="residualRef" class="pred-chart pred-chart-residual"></div> <!-- 预测数据表格 --> <el-table :data="tableData" stripe border style="width:100%"> <el-table-column prop="window_time" label="时间窗口" min-width="170"> <template #default="{ row }">{{ formatWindowTime(row.window_time) }}</template> </el-table-column> <el-table-column prop="true_sales" label="真实销售额" min-width="130"> <template #default="{ row }"> <span style="color:#409eff;font-weight:600">¥{{ row.true_sales }}</span> </template> </el-table-column> <el-table-column prop="pred_sales" label="预测销售额" min-width="130"> <template #default="{ row }"> <span style="color:#67c23a;font-weight:600">¥{{ row.pred_sales }}</span> </template> </el-table-column> <el-table-column label="误差" min-width="120"> <template #default="{ row }"> <span :style="{ color: Math.abs(row.true_sales - row.pred_sales) > 500 ? '#f56c6c' : '#909399' }"> ¥{{ (row.true_sales - row.pred_sales).toFixed(2) }} </span> </template> </el-table-column> <el-table-column prop="create_time" label="生成时间" min-width="170"> <template #default="{ row }">{{ formatDateTime(row.create_time) }}</template> </el-table-column> </el-table> <el-pagination style="margin-top:16px;justify-content:flex-end" v-model:current-page="page" v-model:page-size="size" :total="total" layout="total, prev, pager, next" @change="loadTable" /> </div> </div> </template> <script setup> import { ref, onMounted, onUnmounted } from 'vue' import * as echarts from 'echarts' import request from '@/utils/request' import { formatDateTime, formatWindowTime } from '@/utils/format' const errorMetric = ref({ rmse: 0, mae: 0, mape: 0 }) const tableData = ref([]) const page = ref(1) const size = ref(10) const total = ref(0) const compareRef = ref(null) const residualRef = ref(null) let charts = [] /** * X 轴日期时间标签配置(分行显示,避免底部裁切) */ function buildAxisLabel() { return { rotate: 30, interval: 'auto', hideOverlap: true, fontSize: 11, margin: 16, formatter(val) { const text = formatWindowTime(val) if (text.length >= 16) return `${text.slice(0, 10)}\n${text.slice(11)}` return text }, } } /** * 初始化真实销售额 vs 预测销售额对比图 */ function initCompareChart(data) { const chart = echarts.init(compareRef.value) const labels = data.map(d => formatWindowTime(d.window_time)) chart.setOption({ title: { text: '真实销售额 vs 预测销售额 对比', left: 'center', textStyle: { fontSize: 15 } }, tooltip: { trigger: 'axis', formatter(params) { const idx = params[0]?.dataIndex ?? 0 const lines = [labels[idx] || ''] params.forEach(p => lines.push(`${p.marker}${p.seriesName}: ${p.value}`)) return lines.join('<br/>') }, }, // 图例放顶部,避免与底部日期重叠 legend: { data: ['真实销售额', '预测销售额'], top: 32 }, xAxis: { type: 'category', data: labels, axisTick: { alignWithLabel: true }, axisLabel: buildAxisLabel(), }, yAxis: { type: 'value', name: '销售额(元)' }, series: [ { name: '真实销售额', type: 'line', smooth: true, data: data.map(d => Number(d.true_sales)), itemStyle: { color: '#409eff' }, lineStyle: { width: 3 }, symbol: 'circle', symbolSize: 8, }, { name: '预测销售额', type: 'line', smooth: true, data: data.map(d => Number(d.pred_sales)), itemStyle: { color: '#67c23a' }, lineStyle: { width: 3, type: 'dashed' }, symbol: 'diamond', symbolSize: 8, }, ], grid: { left: 20, right: 24, bottom: 28, top: 72, containLabel: true }, }) charts.push(chart) } /** * 初始化预测残差分析图 */ function initResidualChart(data) { const chart = echarts.init(residualRef.value) const labels = data.map(d => formatWindowTime(d.window_time)) chart.setOption({ title: { text: '预测残差分析 (真实值 - 预测值)', left: 'center', textStyle: { fontSize: 15 } }, tooltip: { trigger: 'axis', formatter(params) { const idx = params[0]?.dataIndex ?? 0 const p = params[0] return `${labels[idx] || ''}<br/>${p.marker}残差: ${p.value}` }, }, xAxis: { type: 'category', data: labels, axisTick: { alignWithLabel: true }, axisLabel: buildAxisLabel(), }, yAxis: { type: 'value', name: '残差(元)' }, series: [{ type: 'bar', data: data.map(d => ({ value: d.residual, itemStyle: { color: d.residual >= 0 ? '#409eff' : '#f56c6c' }, })), barWidth: 20, }], grid: { left: 20, right: 24, bottom: 28, top: 56, containLabel: true }, }) charts.push(chart) } /** * 加载预测图表与误差指标 */ async function loadData() { const [errorRes, compareRes, residualRes] = await Promise.all([ request.get('/prediction/error'), request.get('/prediction/compare'), request.get('/prediction/residual'), ]) errorMetric.value = errorRes.data charts.forEach(c => c.dispose()) charts = [] initCompareChart(compareRes.data) initResidualChart(residualRes.data) } /** * 分页加载预测结果表格 */ async function loadTable() { const res = await request.get('/prediction/list', { params: { page: page.value, size: size.value } }) tableData.value = res.data.items total.value = res.data.total } onMounted(() => { loadData(); loadTable() }) onUnmounted(() => charts.forEach(c => c.dispose())) </script> <style scoped> /* 预留足够高度,保证倾斜日期时间不被裁切 */ .pred-chart { width: 100%; margin-bottom: 24px; } .pred-chart-compare { height: 480px; } .pred-chart-residual { height: 420px; } </style>