
简介面向计算机毕设与课程设计的Spark外卖大数据分析平台项目包围绕外卖订单、用户行为与商家数据的采集、清洗、聚合及可视化展开适合需要完成完整大数据项目或学习Spark开发流程的学生与开发者使用。资源共40个文件核心为14个scala源文件、2个sql脚本与3个hsql脚本另有6个Markdown说明文档、2个json与2个xml配置、1个Python脚本及截图等压缩包仅646KB目录结构清晰便于按模块查阅。已有326人学习下载。包内含Spark程序源码、SQL/Hive分析脚本、pom.xml依赖配置、README说明和项目目录整理可对照理解从数据导入、清洗转换、分析建模到结果导出的完整链路平台涉及HDFS存储、Spark SQL统计、MLlib机器学习与可视化展示等模块分析指标覆盖用户画像、销量预测与商家运营洞察既能作为毕设与课程设计的项目参考也可作为Spark大数据分析入门到实践的案例。1. 为什么外卖分析必须上 Spark一份外卖订单表线上库动辄几亿行叠加骑手位置、商家评分、优惠券核销日志常规关系库跑一次日聚合得等上几分钟。这个《基于Spark的外卖大数据平台分析系统》zip 包把“接数据 → 洗数据 → 算指标 → 出模型 → 上大屏”整条链路用 Spark 串了起来拿到手不是只有几个查询脚本而是带完整 Maven 工程和模块化目录。它适合做大数据方向毕设、课程设计也适合刚从 Hadoop 转 Spark 的工程师当作可落地的参考。对外卖业务来说你需要的不只是“算出结果”而是能在分钟级重新跑完整份全量数据Spark 的内存计算和 DAG 调度正好补上这一环。2. 外卖数据平台的分层架构与 Spark 选型理由2.1 数据分层采集、存储、计算、服务、展示外卖数据从源头到消费我习惯拆成五层。采集层用 Flume 或 Kafka 收集订单日志、用户点击流存储层落 HDFS 做原始备份再按需同步到 HBase 供随机查计算层是 Spark 的主场负责去重、过滤、聚合、训练模型服务层把结果写到 MySQL、Redis 给业务方调用展示层接 ECharts 或 Tableau 做报表。项目里用的就是这套模型各层职责和数据形态如下分层常见组件数据形态核心职责采集Flume / Kafka / Logstash原始日志、JSON、CSV实时或准实时接入订单流存储HDFS / HBaseParquet、Avro、HFile保存全量历史支持成本低的扩容计算Spark Core / SQL / MLlibDataFrame、RDD、Pipeline清洗、聚合、特征工程、训练服务MySQL / Redis结果表、KV供后端 API 和报表查询展示ECharts / Superset / Tableau图表、大屏、仪表盘把指标转化为经营动作这五层不是机械堆叠。对毕设项目来说Kafka 和 HBase 可以先用文件模拟但 Spark 计算层不能省。因为外卖分析的核心矛盾是“数据量大、要求快”只有分布式计算层能让几十亿行数据在数百台机器上并行处理。2.2 为什么存 HDFS 而不是只靠 MySQL外卖订单表有很强的“写一次、读很多次”特征而且分析场景经常要全表扫描。MySQL 的单表千万级之后索引维护和聚合查询会显著变慢。HDFS 的块存储和副本机制更适合大文件顺序读配合 Spark 的本地性调度可以把计算尽量发到数据所在的节点上。项目里的订单历史通常以日期分区存成 Parquet既压缩了存储又用谓词下推省掉无关分区扫描。2.3 pom.xml 依赖与项目代码结构Maven 工程的 pom.xml 是理解 Spark 项目的入口。你打开它一眼能看到依赖作用和版本导向。下面是我在类似项目里复用的一段依赖组合properties spark.version3.1.2/spark.version scala.version2.12/scala.version /properties dependencies dependency groupIdorg.apache.spark/groupId artifactIdspark-core_${scala.version}/artifactId version${spark.version}/version /dependency dependency groupIdorg.apache.spark/groupId artifactIdspark-sql_${scala.version}/artifactId version${spark.version}/version /dependency dependency groupIdorg.apache.spark/groupId artifactIdspark-mllib_${scala.version}/artifactId version${spark.version}/version /dependency /dependencies这段配置说明三件事。第一Spark 3.x 对应的 Scala 2.12依赖的 artifactId 必须带 Scala 版本后缀否则会拉到不兼容的包。第二spark-sql 和 spark-mllib 会传递引入 spark-core不需要重复写。第三如果你用 PySpark则无需 Maven直接pip install pyspark即可。项目里的 main 目录通常按common、etl、analysis、ml分包和后面的数据处理流程一一对应。2.4 本地运行模式和集群模式的切换拿到源码包后建议你先在本地 IDEA 里以local[*]跑通。代码里的SparkSession.builder().master()留了配置位本地开发时设置成local[4]表示用 4 个线程模拟分布式需要跑全量数据再改成 YARN 客户端模式。这种切换对 Spark 程序是透明的也是 Spark 比 MapReduce 易用的点。3. 数据清洗与特征转换从订单日志到分析宽表3.1 订单数据长什么样外卖订单原始日志通常是 CSV 或 JSON常见字段如下。清洗前先建一个字段标准比边写边改更稳。字段类型示例清洗目标order_idLong1029384756主键去重依据user_idLong834201关联用户维度shop_idLong2201关联商家维度amountDouble28.50必须大于 0且为 2 位小数statusStringfinished / cancelled只保留 finishedcreate_timeString2025-06-01 12:30:45转成 Timestamp取小时、星期如果原始文件有几十个字段不要照单全收。只保留分析字段能大幅减少 shuffle 传输量。项目里的订单表经过裁剪后分析用字段不超过 16 个。3.2 用 DataFrame API 做去重和过滤清洗逻辑用 DataFrame API 写会比 RDD 简洁。下面是读取 CSV 并按规则清洗的核心代码val orders spark.read .option(header, true) .option(inferSchema, true) .csv(/data/orders.csv) val ordersClean orders .filter($status finished $amount 0.0) .dropDuplicates(order_id) .select( $order_id, $user_id, $shop_id, $amount, to_timestamp($create_time, yyyy-MM-dd HH:mm:ss).alias(ts) )这段代码先读 CSV 并自动推断类型然后过滤掉取消单和金额异常单再按订单号去重。注意dropDuplicates(order_id)传的字段是去重键如果不传则按整行去重这会导致同一订单因字段微小差异而残留。to_timestamp的第二个参数必须匹配源数据格式如果格式不统一可以先unix_timestamp处理再转换为时间戳。3.3 时间字段解析与维度表关联分析外卖业务离不开时间维度。我的做法是直接从时间字段提取小时、星期、是否周末作为后续模型的输入特征import org.apache.spark.sql.functions._ val ordersWithTime ordersClean .withColumn(hour, hour($ts)) .withColumn(weekday, date_format($ts, E)) // Mon, Tue ... .withColumn(is_weekend, if (weekday is Sunday/Saturday) 1 else 0) val ordersJoined ordersWithTime .join(shopInfo, Seq(shop_id), left_outer)这里hour($ts)使用的是 Spark SQL 内建函数不需要自己写 UDF。date_format($ts, E)返回星期的英文缩写后续可以映射成数值。join时如果两边键都不大默认的SortMergeJoin还好但如果shop_id分布极不均匀就要考虑加盐分桶这个在 3.4 展开。3.4 数据倾斜的加盐处理外卖商家的订单量天然是“二八分布”少数头部商家占了大半订单。直接做 join 或 groupBy 时某个 key 所在的分区会成为长尾。常见的缓解思路是加盐给大 key 追加随机前缀把小 key 放大成多个子键。下面是一个演示扩张商家 ID 的代码思路val saltedOrders ordersJoined .withColumn(salt, floor(rand() * 100)) .withColumn(salted_shop_id, concat($shop_id, lit(_), $salt)) // 对维度表也进行一系列前缀扩展 val shopInfoSalted shopInfo .crossJoin(spark.range(100).withColumnRenamed(id, salt)) .withColumn(salted_shop_id, concat($shop_id, lit(_), $salt))关键点是两边使用相同的盐集合join 才能命中。这个技巧能显著缓解热点 key但不要滥用因为维度表会被放大 100 倍适合少数热点 key 也适合全量放大反而浪费资源。4. Spark SQL 聚合与 MLlib 建模用户画像和销量预测4.1 高峰时段 GMV 统计清洗完的数据可以注册成临时表直接写 SQL。统计各小时订单量与 GMV 的代码SELECT hour, COUNT(*) AS order_cnt, ROUND(SUM(amount), 2) AS gmv FROM orders_clean WHERE ts date_sub(CURRENT_DATE, 30) GROUP BY hour ORDER BY gmv DESCdate_sub(CURRENT_DATE, 30)把统计范围限制在最近 30 天减少无谓扫描。外卖通常出现午晚高峰两个波峰这时候运营活动需要提前安排运力。Spark SQL 的字符串到 DataFrame 的转换只需要spark.sql(...)底层走的是 Catalyst 优化器和手动写 DataFrame 是同一套执行计划。4.2 RFM 用户分层用户分层最经典的模型是 RFM最近一次下单时间Recency、下单频率Frequency、消费金额Monetary。在 Spark 里用 groupBy 一次聚合出来SELECT user_id, DATEDIFF(CURRENT_DATE, MAX(DATE(ts))) AS recency, COUNT(DISTINCT DATE(ts)) AS frequency, SUM(amount) AS monetary FROM orders_clean GROUP BY user_id拿到这三维数值后按分位数切成高/低两档。我给每个维度按 0/1 打分组合成 8 类用户然后用CASE WHEN写进特征表。用户类型RFM运营策略重要价值用户111私域维护、新品推送潜力用户110增加推送频次新用户100发券促活流失风险101召回短信打分阈值我一般用percentile_approx获取 50% 分位点避免精确计算超大字段浪费内存。4.3 随机森林回归预测明日订单量预测外卖门店次日订单量我用 Spark MLlib 的随机森林回归器。特征选历史 7 日订单量、当天星期、是否节假日、门店评分、优惠券力度。特征向量化用VectorAssemblerimport org.apache.spark.ml.feature.VectorAssembler import org.apache.spark.ml.regression.RandomForestRegressor val featureCols Array( avg_7d_orders, avg_7d_amount, weekday_idx, is_holiday, shop_score, coupon_rate ) val assembler new VectorAssembler() .setInputCols(featureCols) .setOutputCol(features) val featureDf assembler.transform(shopFeatures) val Array(train, test) featureDf.randomSplit(Array(0.8, 0.2), seed 42) val rf new RandomForestRegressor() .setLabelCol(next_day_orders) .setFeaturesCol(features) .setNumTrees(80) .setMaxDepth(10) .setMaxBins(32) val model rf.fit(train)setNumTrees(80)控制决策树数量越多越稳但训练越慢setMaxDepth(10)防止单棵树过深带来的过拟合setMaxBins(32)决定连续特征离散化桶数值越大能捕捉的细节越多但内存占用也越高。训练完用测试集计算 RMSEval predictions model.transform(test) val rmse predictions.selectExpr(sqrt(avg(pow(prediction - next_day_orders, 2)))).first().getDouble(0)当 RMSE 稳定在峰值订单量的 10% 以内模型算能用了。如果效果差先回看特征而不是急着调参。4.4 分析结果写回 MySQL 与 HDFS分析结果要给后端展示不能总跑 Spark 任务读 HDFS。我的习惯是把聚合结果和模型预测写回 MySQLval props new java.util.Properties() props.setProperty(user, root) props.setProperty(password, ******) props.setProperty(driver, com.mysql.cj.jdbc.Driver) result .write.mode(overwrite) .jdbc(jdbc:mysql://192.168.1.20:3306/analysis_db, order_stat, props)mode(overwrite)表示全量覆盖适合结果表只有几万行的情况。如果表会越来越大改成append并加日期分区字段。注意driver必须指定否则 Spark 也找不到 MySQL 驱动pom.xml 里要加mysql-connector-java。5. 集群搭建、spark-submit 提交与内存调优5.1 三节点 Spark Standalone 搭建过程拿到毕设源码后很多同学会在本地跑完就结束但面试官更在意你有没有集群部署经验。我之前在三台 8C16G 机器上搭建过流程很固定。先安装 JDK 和 Spark配置环境变量然后修改conf/spark-env.shexport JAVA_HOME/usr/lib/jvm/java-8-openjdk-amd64 export SPARK_MASTER_HOSTnode01 export SPARK_MASTER_PORT7077 export SPARK_WORKER_CORES6 export SPARK_WORKER_MEMORY12g接着在conf/slaves里写入节点列表node01 node02 node03启动后访问node01:8080能看到三节点的 Worker 状态。这里容易踩的坑是机器内存没预留系统余量SPARK_WORKER_MEMORY设置过高导致 OOM。建议每台 16G 内存最多给 Spark 12G留 4G 给 OS 和 HDFS 守护进程。5.2 Spark on YARN 提交命令与参数解析生产环境我一般选 Spark on YARN这样和 HDFS、Hive 共用资源统一调度。提交命令spark-submit \ --master yarn \ --deploy-mode client \ --name order-analysis \ --num-executors 6 \ --executor-cores 3 \ --executor-memory 6g \ --driver-memory 4g \ --class com.example.OrderAnalysis \ order-analysis-1.0.jar参数含义很明确num-executors是启动的 Executor 总数executor-cores和executor-memory是每个 Executor 分到的资源。注意实际的并行度是num-executors * executor-cores如果这个乘积小于spark.sql.shuffle.partitionsshuffle 任务的默认分区数就要相应调低否则会产生大量空任务。实践里executor-memory不能超过 YARN 容器最大值且要留 10% 的 overhead 给 JVM超过会被 ResourceManager 直接杀掉。5.3 内存与 Shuffle 调优参数Spark 内存参数是调优大头外卖这种 groupBy 密集场景尤其重要。下面是我常用的一套起步参数参数起始值适用场景调优依据spark.executor.memory6g通用内存溢出时增加spark.memory.offHeap.enabledfalse大多数情况需要序列化存大对象时开启spark.sql.shuffle.partitions200聚合/join分区数 Executor 核数 × 3 倍spark.default.parallelism100RDD 操作分区数量应大于 Executor 总核数spark.shuffle.file.buffer32k大量小文件增大到 64k 可减少磁盘IOspark.sql.autoBroadcastJoinThreshold10m小表 join维度表小于此阈值自动广播spark.sql.shuffle.partitions是最常见的问题来源。默认 200 在数据量小时偏多在集群规模变大时又偏少。我一般按“每个 Executor 核数 × 3”估算再根据实际某个 stage 的耗时和 GC 情况微调。对 6 个 Executor、每个 3 核来说设 54 比 200 更合适。5.4 分析结果的可视化大屏外卖分析最终要落到业务可看的数据大屏。如果是 React TypeScript 的前端我会把 Spark 写回 MySQL 的结果通过 Node 服务暴露成 JSON前端用 ECharts 渲染。后端返回的指标包括今日 GMV、高峰时段、实时完成率等。如果前端不想写代码也可以用 Superset 直接连 MySQL 做看板。对毕设来说做一个 1920×1080 的大屏展示比单纯截图报表更有说服力。6. 进阶用 Spark Structured Streaming 做外卖实时看板6.1 从批处理复用代码到流处理前面的分析都是批处理订单产生后要延迟一段时间才能知道结果。想看到分钟级的实时数据可以把清洗逻辑抽成函数让批流共用同一套代码。Spark Structured Streaming 的 DataFrame API 与批处理几乎一致只是最后用writeStream输出。6.2 一个可跑的 Kafka 消费聚合例子假设订单消息写到 Kafkatopic 为orders结构是 JSON。实时统计每分钟各店铺的 GMVval stream spark.readStream .format(kafka) .option(kafka.bootstrap.servers, 192.168.1.10:9092) .option(subscribe, orders) .option(startingOffsets, latest) .load() val parsed stream .selectExpr(CAST(value AS STRING) AS json) .selectExpr(get_json_object(json, $.shop_id) AS shop_id, get_json_object(json, $.amount) AS amount, get_json_object(json, $.create_time) AS ts) val windowed parsed .groupBy(window(to_timestamp($ts), 1 minute), $shop_id) .agg(sum($amount.cast(double)).alias(gmv))get_json_object是 Spark SQL 内建的 JSON 解析函数不必引入额外库。window(to_timestamp($ts), 1 minute)定义了一分钟的滑动窗口注意窗口是事件时间不是处理时间。如果消息延迟比较多需要调watermark否则迟到的数据会被丢弃。最后启动写入 HDFS 或控制台windowed.writeStream .outputMode(append) // 或 complete取决于你要全量结果还是增量 .format(console) .option(truncate, false) .trigger(org.apache.spark.sql.streaming.Trigger.ProcessingTime(1 minute)) .start()trigger控制微批频率StartingOffsets设置latest表示只消费新数据。调试时用console输出确认无误再改成foreachBatch写 MySQL。这样整个外卖平台就同时具备了批量和实时的分析能力而 Spark 的 API 从批到流迁移成本非常低。把批处理逻辑抽成def process(orders: DataFrame): DataFrame流里面直接调用同一个函数能做到一份业务代码服务两个场景。本文还有配套的精品资源点击获取